From 1d404db209e48d8c87500d1757260a408c6ec417 Mon Sep 17 00:00:00 2001 From: "NG (Graham)" Date: Sun, 6 Jul 2025 17:36:52 -0400 Subject: [PATCH] Get part way through loading (to sync) #30 --- rc_chat_room/src/operations/join_channel.rs | 4 +- rc_chat_room/src/operations/more_auth.rs | 2 +- rc_chat_room/src/operations/send_message.rs | 8 +- rc_chat_room/src/state/chat/chat.rs | 12 +- rc_core/src/data/cube_list.rs | 13 + rc_core/src/persist/combat.rs | 16 +- rc_core/src/persist/user/account_json.rs | 52 ++- rc_core/src/persist/user/mod.rs | 2 +- rc_core/src/persist/user/traits.rs | 17 +- rc_core/src/state.rs | 8 +- rc_multiplayer/src/events/activate_sync.rs | 30 ++ .../src/events/all_loading_progress.rs | 30 ++ rc_multiplayer/src/events/loading_progress.rs | 33 ++ rc_multiplayer/src/events/mod.rs | 24 +- .../src/events/validate_game_guid.rs | 29 +- rc_multiplayer/src/events/weapon_select.rs | 36 ++ rc_multiplayer/src/handler.rs | 113 +++-- rc_multiplayer/src/handlers/dataless.rs | 29 ++ rc_multiplayer/src/handlers/mod.rs | 1 + rc_multiplayer/src/handlers/simple_typed.rs | 35 +- rc_multiplayer/src/main.rs | 9 +- rc_multiplayer/src/matches/aggregate.rs | 72 ++++ rc_multiplayer/src/matches/engine.rs | 3 + rc_multiplayer/src/matches/generic.rs | 402 ++++++++++++++++++ rc_multiplayer/src/matches/messages.rs | 65 +++ rc_multiplayer/src/matches/mod.rs | 13 + rc_multiplayer/src/traits.rs | 7 +- rc_multiplayer/src/user.rs | 52 ++- rc_multiplayer/src/vehicle_motion.rs | 29 ++ rc_services_room/src/operations/mod.rs | 1 + .../src/operations/singleplayer_campaigns.rs | 2 +- 31 files changed, 1049 insertions(+), 100 deletions(-) create mode 100644 rc_multiplayer/src/events/activate_sync.rs create mode 100644 rc_multiplayer/src/events/all_loading_progress.rs create mode 100644 rc_multiplayer/src/events/loading_progress.rs create mode 100644 rc_multiplayer/src/events/weapon_select.rs create mode 100644 rc_multiplayer/src/handlers/dataless.rs create mode 100644 rc_multiplayer/src/matches/aggregate.rs create mode 100644 rc_multiplayer/src/matches/engine.rs create mode 100644 rc_multiplayer/src/matches/generic.rs create mode 100644 rc_multiplayer/src/matches/messages.rs create mode 100644 rc_multiplayer/src/matches/mod.rs create mode 100644 rc_multiplayer/src/vehicle_motion.rs diff --git a/rc_chat_room/src/operations/join_channel.rs b/rc_chat_room/src/operations/join_channel.rs index b8b9f42..3d9e092 100644 --- a/rc_chat_room/src/operations/join_channel.rs +++ b/rc_chat_room/src/operations/join_channel.rs @@ -25,7 +25,7 @@ async fn do_join_handling(params: ParameterTable<()>, user: &crate::UserTy, chat } let user_info = user.user()?; //let chat_user = super::get_chat_user(user_info.as_ref().as_ref()); - chat_system.system_mut().join_channel(user_info.token().uuid.clone(), chann_name.string.clone()); + chat_system.system_mut().join_channel(user_info.public_id().to_owned(), chann_name.string.clone()); let response = user_info.add_subscribed_channel(chann_name.string, crate::data::channel::ChatChannelType::from_u8(chann_ty as _)?).await?; params.insert(CHANNEL_INFO_PARAM_KEY, response); } @@ -62,7 +62,7 @@ async fn do_leave_handling(params: ParameterTable<()>, user: &crate::UserTy, cha if let Some(Typed::Int(chann_ty)) = params.remove(&CHANNEL_TYPE_PARAM_KEY) { let user_info = user.user()?; //let chat_user = super::get_chat_user(user_info.as_ref().as_ref()); - chat_system.system_mut().leave_channel(user_info.token().uuid.clone(), chann_name.string.clone()); + chat_system.system_mut().leave_channel(user_info.public_id().to_owned(), chann_name.string.clone()); user_info.remove_subscribed_channel(chann_name.string, crate::data::channel::ChatChannelType::from_u8(chann_ty as _)?).await?; } } diff --git a/rc_chat_room/src/operations/more_auth.rs b/rc_chat_room/src/operations/more_auth.rs index fd843e5..5f563f3 100644 --- a/rc_chat_room/src/operations/more_auth.rs +++ b/rc_chat_room/src/operations/more_auth.rs @@ -34,7 +34,7 @@ impl MoreLobbyAuth { if let Some(Typed::Str(auth_payload)) = params.get(&Self::AUTH_PAYLOAD_KEY) { if user.update_with_auth(&auth_payload.string).await { let user_impl = user.user()?; - let name = user_impl.token().uuid.clone(); + let name = user_impl.public_id().to_owned(); //let chat_user = super::get_chat_user(user_impl.as_ref().as_ref()); let channels = user_impl.subscribed_channels_strings().await?; let event_tx = user.event_chann(); diff --git a/rc_chat_room/src/operations/send_message.rs b/rc_chat_room/src/operations/send_message.rs index a4d8335..55c9d49 100644 --- a/rc_chat_room/src/operations/send_message.rs +++ b/rc_chat_room/src/operations/send_message.rs @@ -16,7 +16,7 @@ pub fn send_public_message_handler(chat_system: crate::state::chat::ChatImpl) -> if let Some(Typed::Str(message_text)) = params.remove(&MESSAGE_TEXT_PARAM_KEY) { let user = user.user()?; if message_text.string.bytes().len() > MAX_MESSAGE_LEN { - log::warn!("Rejecting too long chat message from {}", user.token().uuid); + log::warn!("Rejecting too long chat message from {}", user.public_id()); return Err(oj_rc_core::data::error_codes::ChatErrorCodes::Flood as i16) } let chat_loc = if let Some(Typed::Str(chat_loc)) = params.remove(&CHAT_LOCATION_PARAM_KEY) { @@ -25,7 +25,7 @@ pub fn send_public_message_handler(chat_system: crate::state::chat::ChatImpl) -> "".to_owned() }; let chat_system = chat.system(); - log::debug!("Got message `{}` from user {} ({} @ {}/{:?})", message_text.string, user.token().uuid, chat_loc, channel_name.string, channel_enum); + log::debug!("Got message `{}` from user {} ({} @ {}/{:?})", message_text.string, user.public_id(), chat_loc, channel_name.string, channel_enum); chat_system.handle_public_message(user.as_ref().as_ref(), message_text.string, channel_name.string, channel_enum); } } @@ -43,7 +43,7 @@ pub fn send_private_message_handler(chat_system: crate::state::chat::ChatImpl) - if let Some(Typed::Str(message_text)) = params.remove(&MESSAGE_TEXT_PARAM_KEY) { let user = user.user()?; if message_text.string.bytes().len() > MAX_MESSAGE_LEN { - log::warn!("Rejecting too long chat message from {}", user.token().uuid); + log::warn!("Rejecting too long chat message from {}", user.public_id()); return Err(oj_rc_core::data::error_codes::ChatErrorCodes::Flood as i16) } let chat_loc = if let Some(Typed::Str(chat_loc)) = params.remove(&CHAT_LOCATION_PARAM_KEY) { @@ -52,7 +52,7 @@ pub fn send_private_message_handler(chat_system: crate::state::chat::ChatImpl) - "".to_owned() }; let chat_system = chat.system(); - log::debug!("Got message `{}` from user {} (@ {} to {})", message_text.string, user.token().uuid, chat_loc, username.string); + log::debug!("Got message `{}` from user {} (@ {} to {})", message_text.string, user.public_id(), chat_loc, username.string); chat_system.handle_private_message(user.as_ref().as_ref(), message_text.string, username.string); } } diff --git a/rc_chat_room/src/state/chat/chat.rs b/rc_chat_room/src/state/chat/chat.rs index d47a951..874ab46 100644 --- a/rc_chat_room/src/state/chat/chat.rs +++ b/rc_chat_room/src/state/chat/chat.rs @@ -85,13 +85,13 @@ impl ChatSystem { pub fn handle_public_message(&self, user: &dyn oj_rc_core::persist::user::User<()>, text: String, channel: String, channel_ty: crate::data::channel::ChatChannelType) { if self.config.is_command_channel(&channel) { - if let Some(user_handle) = self.online_users.get(&user.token().uuid) { + if let Some(user_handle) = self.online_users.get(user.public_id()) { self.handle_public_command(user, text, user_handle, channel, channel_ty); } } else if let Some(room) = self.chats.get(&channel) { let event_params = crate::events::chat_message::PublicMessage { - sender_name: user.token().uuid.clone(), - sender_display_name: user.token().uuid.clone(), + sender_name: user.public_id().to_owned(), + sender_display_name: user.public_id().to_owned(), text, is_dev: user.is_dev(), is_mod: user.is_mod(), @@ -128,13 +128,13 @@ impl ChatSystem { pub fn handle_private_message(&self, user: &dyn oj_rc_core::persist::user::User<()>, text: String, recipient: String) { if self.config.is_command_user(&recipient) { - if let Some(user_handle) = self.online_users.get(&user.token().uuid) { + if let Some(user_handle) = self.online_users.get(user.public_id()) { self.handle_private_command(user, text, user_handle); } } else if let Some(recipient_handle) = self.online_users.get(&recipient) { let private_msg = crate::events::chat_message::PrivateMessage { - sender_name: user.token().uuid.clone(), - sender_display_name: user.token().uuid.clone(), + sender_name: user.public_id().to_owned(), + sender_display_name: user.public_id().to_owned(), text, is_dev: user.is_dev(), is_mod: user.is_mod(), diff --git a/rc_core/src/data/cube_list.rs b/rc_core/src/data/cube_list.rs index f6a86b7..01d491f 100644 --- a/rc_core/src/data/cube_list.rs +++ b/rc_core/src/data/cube_list.rs @@ -105,6 +105,19 @@ impl ItemTier { ItemTier::T5 => "T5", } } + + pub fn from_u32(num: u32) -> Option { + match num { + 0 => Some(Self::NoTier), + 100 => Some(Self::T0), + 200 => Some(Self::T1), + 300 => Some(Self::T2), + 400 => Some(Self::T3), + 500 => Some(Self::T4), + 600 => Some(Self::T5), + _ => None, + } + } } #[derive(Clone, Copy)] diff --git a/rc_core/src/persist/combat.rs b/rc_core/src/persist/combat.rs index fd4efeb..b380abe 100644 --- a/rc_core/src/persist/combat.rs +++ b/rc_core/src/persist/combat.rs @@ -361,7 +361,7 @@ fn default_rotation() -> GameEventSequence { multiplayer: GameEvent { map: GameMap::Neptune3, visibility: GameVisibility::Poor, - mode: GameType::BattleArena, + mode: GameType::SuddenDeath, auto_heal: true, }, duration_s: 5*60, // 5 minutes @@ -376,7 +376,7 @@ fn default_rotation() -> GameEventSequence { multiplayer: GameEvent { map: GameMap::Mars1, visibility: GameVisibility::Poor, - mode: GameType::Pit, + mode: GameType::SuddenDeath, auto_heal: true, }, duration_s: 5*60, @@ -391,7 +391,7 @@ fn default_rotation() -> GameEventSequence { multiplayer: GameEvent { map: GameMap::Mars1, visibility: GameVisibility::Poor, - mode: GameType::TestMode, + mode: GameType::SuddenDeath, auto_heal: true, }, duration_s: 5*60, @@ -406,7 +406,7 @@ fn default_rotation() -> GameEventSequence { multiplayer: GameEvent { map: GameMap::Neptune3, visibility: GameVisibility::Poor, - mode: GameType::BattleArena, + mode: GameType::SuddenDeath, auto_heal: true, }, duration_s: 5*60, @@ -421,7 +421,7 @@ fn default_rotation() -> GameEventSequence { multiplayer: GameEvent { map: GameMap::Mars1, visibility: GameVisibility::Poor, - mode: GameType::Pit, + mode: GameType::SuddenDeath, auto_heal: true, }, duration_s: 5*60, @@ -436,7 +436,7 @@ fn default_rotation() -> GameEventSequence { multiplayer: GameEvent { map: GameMap::Mars1, visibility: GameVisibility::Poor, - mode: GameType::TestMode, + mode: GameType::SuddenDeath, auto_heal: true, }, duration_s: 5*60, @@ -451,7 +451,7 @@ fn default_rotation() -> GameEventSequence { multiplayer: GameEvent { map: GameMap::Mars1, visibility: GameVisibility::Poor, - mode: GameType::Pit, + mode: GameType::SuddenDeath, auto_heal: true, }, duration_s: 5*60, @@ -466,7 +466,7 @@ fn default_rotation() -> GameEventSequence { multiplayer: GameEvent { map: GameMap::Mars1, visibility: GameVisibility::Poor, - mode: GameType::TestMode, + mode: GameType::SuddenDeath, auto_heal: true, }, duration_s: 5*60, diff --git a/rc_core/src/persist/user/account_json.rs b/rc_core/src/persist/user/account_json.rs index 94dfb80..b320e42 100644 --- a/rc_core/src/persist/user/account_json.rs +++ b/rc_core/src/persist/user/account_json.rs @@ -59,7 +59,7 @@ impl AccountProvider { #[async_trait::async_trait] impl super::UserProvider for AccountProvider { - async fn authenticate(&self, token: super::UserToken, ext: std::collections::HashMap>) -> Result + Send + Sync>, String> { + async fn authenticate(&self, token: super::UserToken) -> Result + Send + Sync>, String> { //let new_root = self.root.join(&token.uuid); let secret = jsonwebtoken::DecodingKey::from_secret(&self.secret); let mut validation = jsonwebtoken::Validation::new(jsonwebtoken::Algorithm::HS256); @@ -77,16 +77,34 @@ impl super::UserProvider for AccountProvider { }; //let account_info = AccountInfo::load(&new_root).map_err(|e| e.to_string())?; Ok(Box::new(UserData { - token, account: user_info, perms: user_perms, cubes: self.cubes.clone(), garage_upgrades: self.garage_upgrades.clone(), - extensions: ext, db: self.db.clone(), })) //Err("Unable to authenticate".to_string()) } + + async fn multiplayer_authenticate(&self, user: String) -> Result + Send + Sync>, String> { + let user_info = if let Some(user_info) = self.db.user_by_display_name(user).await.map_err(|e| e.to_string())? { + user_info + } else { + return Err("User not found".to_owned()); + }; + let user_perms = if let Some(user_perms) = self.db.perms_by_user_id(user_info.id).await.map_err(|e| e.to_string())? { + user_perms + } else { + return Err("User permissions not found".to_owned()); + }; + Ok(Box::new(UserData { + account: user_info, + perms: user_perms, + cubes: self.cubes.clone(), + garage_upgrades: self.garage_upgrades.clone(), + db: self.db.clone(), + })) + } } #[async_trait::async_trait] @@ -193,12 +211,10 @@ impl super::UserAuthenticator for AccountProvider { #[allow(dead_code)] struct UserData { - token: super::UserToken, account: oj_rc_database::schema::user::Model, perms: oj_rc_database::schema::permissions::Model, cubes: std::sync::Arc>, garage_upgrades: std::sync::Arc, - extensions: std::collections::HashMap>, db: std::sync::Arc, } @@ -270,7 +286,7 @@ impl UserData { log::error!("Failed to find selected vehicle for user_id {} (user_player_data)", self.account.id); polariton_server::operations::SimpleOpError::with_message(INVALID_ROBOT_ERR, "No selected garage".to_owned()) })?; - let user_uuid = self.token.uuid.clone(); + let user_uuid = self.account.public_id.clone(); let weapon_order = oj_rc_database::schema::parse_int_csv(¤t_slot.weapon_order).into_iter().map(|x| x as i32).collect::>(); let user_avatar_aux = self.db.user_aux_by_user_id_and_descriptor(self.account.id, oj_rc_database::schema::user_aux::Descriptor::AvatarId).await.map_err(|e| { log::error!("Failed to retrieve avatar for user_id {} (user_player_data): {}", self.account.id, e); @@ -442,12 +458,8 @@ const UNEXPECTED_ERR: i16 = crate::data::error_codes::WebServicesError::Unexpect #[async_trait::async_trait] impl super::User for UserData { - fn ext(&self, ty: std::any::TypeId) -> Option<&'_ (dyn std::any::Any + Send + Sync + 'static)> { - self.extensions.get(&ty).map(|x| x.as_ref()) - } - - fn token(&self) -> &'_ super::UserToken { - &self.token + fn public_id(&self) -> &'_ str { + &self.account.public_id } fn is_mod(&self) -> bool { @@ -1108,3 +1120,19 @@ impl super::LobbyUser for UserData { }) } } + +#[async_trait::async_trait] +impl super::MultiplayerUser for UserData { + // TODO + fn user_id(&self) -> i32 { + self.account.id + } + + fn user_name(&self) -> &'_ str { + &self.account.public_id + } + + fn display_name(&self) -> &'_ str { + &self.account.display_name + } +} diff --git a/rc_core/src/persist/user/mod.rs b/rc_core/src/persist/user/mod.rs index b857413..c87bb14 100644 --- a/rc_core/src/persist/user/mod.rs +++ b/rc_core/src/persist/user/mod.rs @@ -11,7 +11,7 @@ mod inventory; pub use inventory::UnlockedParts; mod traits; -pub use traits::{UserProvider, User, UserToken, UserSlots, UserSlotData, VehicleData, UserInfo, UserLoginInfo, ExtraUserInfo, UserAuthenticator, NewSlotData, UserId, RegistrationInfo, VehicleUploadData, ChatUser, AvatarInfo, GetAvatarInfo, ControlData, ControlType, CustomisationData, GetCustomisationData, SetSanction, SanctionType, LobbyUser}; +pub use traits::{UserProvider, User, UserToken, UserSlots, UserSlotData, VehicleData, UserInfo, UserLoginInfo, ExtraUserInfo, UserAuthenticator, NewSlotData, UserId, RegistrationInfo, VehicleUploadData, ChatUser, AvatarInfo, GetAvatarInfo, ControlData, ControlType, CustomisationData, GetCustomisationData, SetSanction, SanctionType, LobbyUser, MultiplayerUser}; pub const TOKEN_SECRET_FILENAME: &str = "token_secret.key"; diff --git a/rc_core/src/persist/user/traits.rs b/rc_core/src/persist/user/traits.rs index ea429c8..3142a6c 100644 --- a/rc_core/src/persist/user/traits.rs +++ b/rc_core/src/persist/user/traits.rs @@ -43,7 +43,9 @@ pub struct RegistrationInfo { #[async_trait::async_trait] pub trait UserProvider { - async fn authenticate(&self, user: UserToken, ext: std::collections::HashMap>) -> Result + Send + Sync>, String>; + async fn authenticate(&self, user: UserToken) -> Result + Send + Sync>, String>; + + async fn multiplayer_authenticate(&self, user: String) -> Result + Send + Sync>, String>; } #[async_trait::async_trait] @@ -54,9 +56,8 @@ pub trait UserAuthenticator { } #[async_trait::async_trait] -pub trait User: ChatUser + LobbyUser { - fn ext(&self, ty: std::any::TypeId) -> Option<&'_ (dyn std::any::Any + Send + Sync + 'static)>; - fn token(&self) -> &'_ super::UserToken; +pub trait User: ChatUser + LobbyUser + MultiplayerUser { + fn public_id(&self) -> &'_ str; fn is_mod(&self) -> bool; fn is_admin(&self) -> bool; fn is_dev(&self) -> bool; @@ -230,3 +231,11 @@ impl SanctionType { pub trait LobbyUser { async fn player_data(&self) -> Result; } + +#[async_trait::async_trait] +pub trait MultiplayerUser { + // TODO + fn user_id(&self) -> i32; + fn user_name(&self) -> &'_ str; + fn display_name(&self) -> &'_ str; +} diff --git a/rc_core/src/state.rs b/rc_core/src/state.rs index e983a04..8fb18b7 100644 --- a/rc_core/src/state.rs +++ b/rc_core/src/state.rs @@ -27,12 +27,12 @@ impl UserState { token: splits[1].to_owned(), refresh_token: splits[2].to_owned(), }; - let ext = if let Some(ext) = ext_f(&token) { - ext + if let Some(ext) = ext_f(&token) { + log::debug!("Ignoring extra auth info: {:?}", ext); } else { return false; - }; - match auth.authenticate(token, ext).await { + } + match auth.authenticate(token).await { Ok(user) => { let mut lock = self.state.write().unwrap(); *lock = InitState::Authenticated(std::sync::Arc::new(user)); diff --git a/rc_multiplayer/src/events/activate_sync.rs b/rc_multiplayer/src/events/activate_sync.rs new file mode 100644 index 0000000..68e4694 --- /dev/null +++ b/rc_multiplayer/src/events/activate_sync.rs @@ -0,0 +1,30 @@ +pub struct RequestLoadingSync { + msg_router: tokio::sync::mpsc::Sender, +} + +pub(super) fn handler(init_ctx: &crate::InitConfig) -> crate::handlers::dataless::Dataless { + crate::handlers::dataless::Dataless::new(RequestLoadingSync::new(init_ctx)) +} + +impl RequestLoadingSync { + fn new(init_ctx: &crate::InitConfig) -> Self { + Self { + msg_router: init_ctx.matches_chann.clone(), + } + } +} + +#[async_trait::async_trait] +impl crate::handlers::dataless::DatalessEventCodeHandler for RequestLoadingSync { + const CODE: rlnl::event_code::NetworkEvent = rlnl::event_code::NetworkEvent::RequestSync; + + async fn handle(&self, _peer: &std::sync::Arc>, user: &crate::UserData, _sender: &std::sync::Arc>) { + if let Some(user_info) = user.user().await { + super::log_channel_send_failure(self.msg_router.send(crate::matches::GameMessage::RequestLoadingSync { + user_id: user_info.user_id(), + }).await); + } else { + log::error!("Failed to handle sync loading request for unknown user"); + } + } +} diff --git a/rc_multiplayer/src/events/all_loading_progress.rs b/rc_multiplayer/src/events/all_loading_progress.rs new file mode 100644 index 0000000..ce2e91c --- /dev/null +++ b/rc_multiplayer/src/events/all_loading_progress.rs @@ -0,0 +1,30 @@ +pub struct RequestAllLoadingProgress { + msg_router: tokio::sync::mpsc::Sender, +} + +pub(super) fn handler(init_ctx: &crate::InitConfig) -> crate::handlers::dataless::Dataless { + crate::handlers::dataless::Dataless::new(RequestAllLoadingProgress::new(init_ctx)) +} + +impl RequestAllLoadingProgress { + fn new(init_ctx: &crate::InitConfig) -> Self { + Self { + msg_router: init_ctx.matches_chann.clone(), + } + } +} + +#[async_trait::async_trait] +impl crate::handlers::dataless::DatalessEventCodeHandler for RequestAllLoadingProgress { + const CODE: rlnl::event_code::NetworkEvent = rlnl::event_code::NetworkEvent::RequestLoadingProgressAllUsers; + + async fn handle(&self, _peer: &std::sync::Arc>, user: &crate::UserData, _sender: &std::sync::Arc>) { + if let Some(user_info) = user.user().await { + super::log_channel_send_failure(self.msg_router.send(crate::matches::GameMessage::RequestLoadingProgress { + user_id: user_info.user_id(), + }).await); + } else { + log::error!("Failed to broadcast loading progress for unknown user"); + } + } +} diff --git a/rc_multiplayer/src/events/loading_progress.rs b/rc_multiplayer/src/events/loading_progress.rs new file mode 100644 index 0000000..df0d670 --- /dev/null +++ b/rc_multiplayer/src/events/loading_progress.rs @@ -0,0 +1,33 @@ +pub struct GameLoadingProgress { + msg_router: tokio::sync::mpsc::Sender, +} + +pub(super) fn handler(init_ctx: &crate::InitConfig) -> crate::handlers::simple_typed::SimpleRlnl { + crate::handlers::simple_typed::SimpleRlnl::new(GameLoadingProgress::new(init_ctx)) +} + +impl GameLoadingProgress { + fn new(init_ctx: &crate::InitConfig) -> Self { + Self { + msg_router: init_ctx.matches_chann.clone(), + } + } +} + +#[async_trait::async_trait] +impl crate::handlers::simple_typed::RlnlEventCodeHandler for GameLoadingProgress { + type In = rlnl::events::loading::LoadingProgress; + const CODE: rlnl::event_code::NetworkEvent = rlnl::event_code::NetworkEvent::BroadcastLoadingProgress; + + async fn handle(&self, data: Self::In, _peer: &std::sync::Arc>, user: &crate::UserData, _sender: &std::sync::Arc>) { + if let Some(user_info) = user.user().await { + super::log_channel_send_failure(self.msg_router.send(crate::matches::GameMessage::LoadingProgress { + user_id: user_info.user_id(), + user_name: data.user_name.0, + progress: data.progress, + }).await); + } else { + log::error!("Failed to broadcast loading progress for user {} (no auth!)", data.user_name.0); + } + } +} diff --git a/rc_multiplayer/src/events/mod.rs b/rc_multiplayer/src/events/mod.rs index 25d45e6..9b79751 100644 --- a/rc_multiplayer/src/events/mod.rs +++ b/rc_multiplayer/src/events/mod.rs @@ -1,6 +1,28 @@ mod validate_game_guid; +mod loading_progress; +mod all_loading_progress; +mod weapon_select; +mod activate_sync; pub async fn handler(init_ctx: &crate::InitConfig) -> crate::handler::LnlEventHandler { - crate::handler::LnlEventHandler::new() + crate::handler::LnlEventHandler::new(init_ctx.users.clone(), crate::vehicle_motion::handler(init_ctx)) .add(validate_game_guid::handler(init_ctx)) + .add(loading_progress::handler(init_ctx)) + .add(all_loading_progress::handler(init_ctx)) + .add(weapon_select::handler(init_ctx)) + .add(activate_sync::handler(init_ctx)) +} + +#[inline] +pub fn log_channel_send_failure(result: Result<(), tokio::sync::mpsc::error::SendError>) { + if result.is_err() { + log::error!("Failed to send game message"); + } +} + +#[inline] +pub fn log_lnl_send_failure(result: std::io::Result) { + if let Err(e) = result { + log::error!("Failed to send packet: {}", e); + } } diff --git a/rc_multiplayer/src/events/validate_game_guid.rs b/rc_multiplayer/src/events/validate_game_guid.rs index f73ab13..41f934f 100644 --- a/rc_multiplayer/src/events/validate_game_guid.rs +++ b/rc_multiplayer/src/events/validate_game_guid.rs @@ -1,5 +1,5 @@ pub struct AuthUserGame { - + matches: tokio::sync::mpsc::Sender, } pub(super) fn handler(init_ctx: &crate::InitConfig) -> crate::handlers::simple_typed::SimpleRlnl { @@ -7,9 +7,9 @@ pub(super) fn handler(init_ctx: &crate::InitConfig) -> crate::handlers::simple_t } impl AuthUserGame { - fn new(_init_ctx: &crate::InitConfig) -> Self { + fn new(init_ctx: &crate::InitConfig) -> Self { Self { - + matches: init_ctx.matches_chann.clone(), } } } @@ -19,7 +19,26 @@ impl crate::handlers::simple_typed::RlnlEventCodeHandler for AuthUserGame { type In = rlnl::events::loading::GameGuidInfo; const CODE: rlnl::event_code::NetworkEvent = rlnl::event_code::NetworkEvent::ValidateGameGuid; - async fn handle(&self, data: Self::In, peer: &std::sync::Arc>, user: &crate::UserData, sender: &literustlib_server::DataSender) { - log::debug!("Got {:?} event with data {:?}", Self::CODE, data); + async fn handle(&self, data: Self::In, peer: &std::sync::Arc>, user: &crate::UserData, sender: &std::sync::Arc>) { + let username = data.player_name.0.clone(); + let game_guid = data.game_guid.0.clone(); + if user.authenticate(data).await { + let user_info = user.user().await.unwrap(); + let (tx, rx) = tokio::sync::oneshot::channel(); + super::log_channel_send_failure(self.matches.send(crate::matches::GameMessage::NewConnection { + user: user_info.clone(), + game_guid, + connection: peer.to_owned(), + response: tx, + sender: sender.to_owned(), + }).await); + log::debug!("Sent NewConnection message to matches handler"); + if let Ok(Some(e)) = rx.await { + log::error!("Failed {:?} event: {}", Self::CODE, e); + } + } else { + log::error!("Failed to validate game guid for user {} (other packets will probably be ignored)", username); + } + } } diff --git a/rc_multiplayer/src/events/weapon_select.rs b/rc_multiplayer/src/events/weapon_select.rs new file mode 100644 index 0000000..a4fb510 --- /dev/null +++ b/rc_multiplayer/src/events/weapon_select.rs @@ -0,0 +1,36 @@ +pub struct WeaponSelect { + msg_router: tokio::sync::mpsc::Sender, +} + +pub(super) fn handler(init_ctx: &crate::InitConfig) -> crate::handlers::simple_typed::SimpleRlnl { + crate::handlers::simple_typed::SimpleRlnl::new(WeaponSelect::new(init_ctx)) +} + +impl WeaponSelect { + fn new(init_ctx: &crate::InitConfig) -> Self { + Self { + msg_router: init_ctx.matches_chann.clone(), + } + } +} + +#[async_trait::async_trait] +impl crate::handlers::simple_typed::RlnlEventCodeHandler for WeaponSelect { + type In = rlnl::events::ingame::SelectWeapon; + const CODE: rlnl::event_code::NetworkEvent = rlnl::event_code::NetworkEvent::WeaponSelect; + + async fn handle(&self, data: Self::In, _peer: &std::sync::Arc>, user: &crate::UserData, _sender: &std::sync::Arc>) { + if let Some(user_info) = user.user().await { + if let Some(category) = oj_rc_core::data::weapon_list::ItemCategory::from_smaller(data.item_category as _) { + if let Some(tier) = oj_rc_core::data::cube_list::ItemTier::from_u32(data.item_category as _) { + super::log_channel_send_failure(self.msg_router.send(crate::matches::GameMessage::WeaponSelect { + user_id: user_info.user_id(), + machine_id: data.machine_id, + category, + size: tier, + }).await); + } else { log::warn!("Bad WeaponSelect tier") } + } else { log::warn!("Bad WeaponSelect category") } + } + } +} diff --git a/rc_multiplayer/src/handler.rs b/rc_multiplayer/src/handler.rs index 5f1ed9f..4d579ba 100644 --- a/rc_multiplayer/src/handler.rs +++ b/rc_multiplayer/src/handler.rs @@ -1,16 +1,22 @@ pub struct LnlEventHandler { event_handlers: std::collections::HashMap>, + motion_handler: Box, + user_provider: std::sync::Arc } impl LnlEventHandler { - pub fn new() -> Self { + pub fn new(user_provider: std::sync::Arc, motion_handler: M) -> Self { Self { event_handlers: std::collections::HashMap::new(), + motion_handler: Box::new(motion_handler), + user_provider, } } pub fn add(mut self, handler: H) -> Self { - self.event_handlers.insert(H::CODE, Box::new(handler)); + if self.event_handlers.insert(H::CODE, Box::new(handler)).is_some() { + log::warn!("Replaced event handler {} with new handler", H::CODE); + } self } } @@ -20,20 +26,33 @@ impl literustlib_server::EventHandler for LnlEventHandler { type PacketData = super::PacketData; type UserData = super::UserData; - async fn on_receive(&self, data: Self::PacketData, _header: &literustlib::packet::Header, peer: &std::sync::Arc< literustlib_server::Connection>, user: &Self::UserData, sender: &literustlib_server::DataSender) { - log::debug!("Got event {:?} (len: {}) from connection id {}", data.variant, data.data.len(), peer.id()); - if let Some(handler) = self.event_handlers.get(&(data.variant as i16)) { - handler.handle(&data.data, peer, user, sender).await; - } else { - #[cfg(debug_assertions)] - { - panic!("Unsupported event variant {:?} ({}), pls fix!!!\n {:?}", data.variant, data.variant as u16, &data.data[..]); - } - #[cfg(not(debug_assertions))] - { - log::warn!("Unsupported event variant {:?} ({}), pls fix!!!", data.variant, data.variant as u16); - } + async fn on_receive(&self, data: Self::PacketData, _header: &literustlib::packet::Header, peer: &std::sync::Arc< literustlib_server::Connection>, user: &Self::UserData, sender: &std::sync::Arc>) { + log::debug!("Got message {:?} (len: {}) from connection id {}", data.message_ty, data.data.len(), peer.id()); + match data.message_ty { + crate::data::MessageType::ClientMsg => { + if let Some(handler) = self.event_handlers.get(&data.variant) { + handler.handle(&data.data, peer, user, sender).await; + } else { + let variant_pretty = i16_to_event(data.variant).map(|x| format!("{:?}", x)).unwrap_or_else(|| "???".to_owned()); + #[cfg(debug_assertions)] + { + panic!("Unsupported event variant {} ({}), pls fix!!!\n {:?}", variant_pretty, data.variant, &data.data[..]); + } + #[cfg(not(debug_assertions))] + { + log::warn!("Unsupported event variant {} ({}), pls fix!!!", variant_pretty, data.variant); + } + } + }, + crate::data::MessageType::ServerMsg => { + log::debug!("Got message from server but I'm the server??? (ignoring)"); + }, + crate::data::MessageType::RobotMotion => { + self.motion_handler.handle(&data.data, user).await; + //log::warn!("Ignoring robot motion message"); + }, } + } async fn on_connect_start(&self, addr: &core::net::SocketAddr, key: String, peer: &std::sync::Arc< literustlib_server::Connection>) -> Option { @@ -41,15 +60,16 @@ impl literustlib_server::EventHandler for LnlEventHandler { //let mut buf = Vec::new(); //literustlib::packet::Packet::with_data(literustlib::packet::Property::Reliable, &[9, 0, 0, 0, 0, 0]).dump(&mut buf).unwrap_or_default(); //socket.send_to(&buf, addr).await.unwrap_or_default(); - Some(crate::UserData::new()) + Some(crate::UserData::new(self.user_provider.clone())) } - async fn on_connect_done(&self, peer: &std::sync::Arc< literustlib_server::Connection>, _user: &Self::UserData, sender: &literustlib_server::DataSender) { + async fn on_connect_done(&self, peer: &std::sync::Arc< literustlib_server::Connection>, _user: &Self::UserData, sender: &std::sync::Arc>) { log::debug!("New connection completed (id:{})", peer.id()); - //let mut buf = Vec::new(); - //literustlib::packet::Packet::with_data(literustlib::packet::Property::Reliable, &[9, 0, 0, 0, 0, 0]).dump(&mut buf).unwrap_or_default(); - //socket.send_to(&buf, addr).await.unwrap_or_default(); - if let Err(e) = sender.send_to(bytes::Bytes::from_static(&[49, 0, 9, 0, 0, 0]), literustlib::packet::Property::Reliable, peer).await { + let data = EventData::without_data( + crate::data::MessageType::ServerMsg, + rlnl::event_code::NetworkEvent::OnConnectedToGameServer, + ); + if let Err(e) = sender.send_data(data, literustlib::packet::Property::Reliable, peer).await { log::error!("Failed to send rlnl OnConnectedToGameServer event: {}", e); } @@ -59,41 +79,60 @@ impl literustlib_server::EventHandler for LnlEventHandler { #[derive(Debug)] pub struct EventData { pub message_ty: crate::data::MessageType, - pub variant: rlnl::event_code::NetworkEvent, + pub variant: i16, pub data_size: u16, pub data: bytes::Bytes, } +impl EventData { + pub fn with_data(message_ty: crate::data::MessageType, event: rlnl::event_code::NetworkEvent, data: bytes::Bytes) -> Self { + Self { + message_ty, + variant: event as i16, + data_size: data.len().try_into().expect("Event data too large"), + data, + } + } + + pub fn without_data(message_ty: crate::data::MessageType, event: rlnl::event_code::NetworkEvent) -> Self { + Self { + message_ty, + variant: event as i16, + data_size: 0, + data: bytes::Bytes::new(), + } + } +} + impl literustlib::packet::PacketData for EventData { fn parse(bytes: bytes::Bytes, _header: &literustlib::packet::Header) -> std::io::Result { + log::debug!("Got packet data ({}) {:?}", bytes.len(), &bytes[..]); if bytes.len() >= 6 { let data = bytes.slice(6..); let net_message_num = i16::from_le_bytes([bytes[0], bytes[1]]); let net_message_type = crate::data::MessageType::from_i16(net_message_num).ok_or_else(|| std::io::Error::new(std::io::ErrorKind::Unsupported, format!("Unsupported message type {}", net_message_num)))?; - let event_code = i16::from_le_bytes([bytes[2], bytes[3]]); + let variant = i16::from_le_bytes([bytes[2], bytes[3]]); let data_size = u16::from_le_bytes([bytes[4], bytes[5]]); - if let Some(event_variant) = i16_to_event(event_code) { - Ok(Self { - message_ty: net_message_type, - variant: event_variant, - data_size, - data, - }) - } else { - Err(std::io::Error::new(std::io::ErrorKind::InvalidData, "Invalid packet code")) - } + Ok(Self { + message_ty: net_message_type, + variant, + data_size, + data, + }) } else { - Err(std::io::Error::new(std::io::ErrorKind::InvalidData, "Packet data is not long enough")) + Err(std::io::Error::new(std::io::ErrorKind::InvalidData, "Packet data is too short")) } } - fn dump(&self) -> Vec { + fn dump(&self) -> bytes::Bytes { use std::io::Write; let mut buf = Vec::new(); - buf.write_all(&(self.variant as u16).to_be_bytes()).unwrap(); + buf.write_all(&(self.message_ty as i16).to_le_bytes()).unwrap(); + buf.write_all(&(self.variant as i16).to_le_bytes()).unwrap(); + buf.write_all(&(self.data_size as u16).to_le_bytes()).unwrap(); buf.write_all(&self.data).unwrap(); - buf + buf.into() } } diff --git a/rc_multiplayer/src/handlers/dataless.rs b/rc_multiplayer/src/handlers/dataless.rs new file mode 100644 index 0000000..200a0de --- /dev/null +++ b/rc_multiplayer/src/handlers/dataless.rs @@ -0,0 +1,29 @@ +pub struct Dataless { + handler: H, +} + +impl Dataless { + pub fn new(inner: H) -> Self { + Self { + handler: inner, + } + } +} + +#[async_trait::async_trait] +pub trait DatalessEventCodeHandler: Sync + Send { + const CODE: rlnl::event_code::NetworkEvent; + + async fn handle(&self, peer: &std::sync::Arc>, user: &crate::UserData, sender: &std::sync::Arc>); +} + +#[async_trait::async_trait] +impl crate::EventCodeHandler for Dataless { + async fn handle(&self, _data: &bytes::Bytes, peer: &std::sync::Arc>, user: &crate::UserData, sender: &std::sync::Arc>) { + self.handler.handle(peer, user, sender).await; + } +} + +impl crate::EventCode for Dataless { + const CODE: i16 = H::CODE as i16; +} diff --git a/rc_multiplayer/src/handlers/mod.rs b/rc_multiplayer/src/handlers/mod.rs index f0b1ad9..4779f2d 100644 --- a/rc_multiplayer/src/handlers/mod.rs +++ b/rc_multiplayer/src/handlers/mod.rs @@ -1 +1,2 @@ pub mod simple_typed; +pub mod dataless; diff --git a/rc_multiplayer/src/handlers/simple_typed.rs b/rc_multiplayer/src/handlers/simple_typed.rs index f8622f0..a445d00 100644 --- a/rc_multiplayer/src/handlers/simple_typed.rs +++ b/rc_multiplayer/src/handlers/simple_typed.rs @@ -18,14 +18,14 @@ pub trait RlnlEventCodeHandler: Sync + Send { //type Out: byteserde::ser_heap::ByteSerializeHeap; const CODE: rlnl::event_code::NetworkEvent; - async fn handle(&self, data: Self::In, peer: &std::sync::Arc>, user: &crate::UserData, sender: &literustlib_server::DataSender); + async fn handle(&self, data: Self::In, peer: &std::sync::Arc>, user: &crate::UserData, sender: &std::sync::Arc>); } #[async_trait::async_trait] impl , H: RlnlEventCodeHandler> crate::EventCodeHandler for SimpleRlnl { - async fn handle(&self, data: &bytes::Bytes, peer: &std::sync::Arc>, user: &crate::UserData, sender: &literustlib_server::DataSender) { + async fn handle(&self, data: &bytes::Bytes, peer: &std::sync::Arc>, user: &crate::UserData, sender: &std::sync::Arc>) { let mut des = byteserde::des_slice::ByteDeserializerSlice::new(&data); - let rlnl_data = In::byte_deserialize(&mut des).expect("Bad serialization"); + let rlnl_data = In::byte_deserialize(&mut des).expect("Bad deserialization"); self.handler.handle(rlnl_data, peer, user, sender).await; } } @@ -34,4 +34,33 @@ impl , H: RlnlEventCodeHandle const CODE: i16 = H::CODE as i16; } +pub struct RlnlSender<'a> { + sender: &'a literustlib_server::DataSender, +} + +impl <'a> RlnlSender<'a> { + #[inline] + pub fn new(inner: &'a literustlib_server::DataSender) -> Self { + Self { + sender: inner, + } + } + + pub async fn send_data(&self, data: &D, event: rlnl::event_code::NetworkEvent, property: literustlib::packet::Property, conn: &literustlib_server::Connection) -> std::io::Result { + let mut ser = byteserde::ser_heap::ByteSerializerHeap::default(); + data.byte_serialize_heap(&mut ser).map_err(|e| std::io::Error::new(std::io::ErrorKind::Unsupported, e.message))?; + let event_data = crate::handler::EventData::with_data( + crate::data::MessageType::ServerMsg, + event, + bytes::Bytes::copy_from_slice(ser.as_slice()), + ); + self.sender.send_data(event_data, property, conn).await + } + + pub async fn send_empty(&self, event: rlnl::event_code::NetworkEvent, property: literustlib::packet::Property, conn: &literustlib_server::Connection) -> std::io::Result { + let event_data = crate::handler::EventData::without_data(crate::data::MessageType::ServerMsg, event); + self.sender.send_data(event_data, property, conn).await + } +} + diff --git a/rc_multiplayer/src/main.rs b/rc_multiplayer/src/main.rs index a424a09..d2ebe4f 100644 --- a/rc_multiplayer/src/main.rs +++ b/rc_multiplayer/src/main.rs @@ -1,16 +1,19 @@ mod cli; mod handler; mod traits; -pub use traits::{EventCodeHandler, UserData, PacketData, EventCode}; +pub use traits::{EventCodeHandler, UserData, PacketData, EventCode, RobotMotionHandler}; mod data; mod events; mod handlers; mod user; +mod matches; +mod vehicle_motion; pub struct InitConfig { pub config: oj_rc_core::persist::config::ConfigImpl, pub users: std::sync::Arc, pub parsers: oj_rc_core::cubes::CubeParsers, + pub matches_chann: tokio::sync::mpsc::Sender, } #[tokio::main] @@ -22,15 +25,19 @@ async fn main() -> std::io::Result<()> { let config = oj_rc_core::persist::config::ConfigImpl::load(&args.assets).expect("Bad config data"); let users = std::sync::Arc::new(oj_rc_core::persist::user::UserImpl::load(&args.data, &config).await.expect("Bad user data")); let parsers = oj_rc_core::cubes::CubeParsers::new(&config); + let matches = matches::GameMatches::new(); + let matches_chann = matches.spawn(); let init_ctx = InitConfig { config, users, parsers, + matches_chann, }; let mtu = oj_rc_core::ConfigProvider::<()>::network_config(&init_ctx.config).max_packet_size; let event_handler = events::handler(&init_ctx).await; let server = literustlib_server::Server::new(event_handler, (args.ip, args.port), mtu).await.expect("Bad server"); + server.listen().await } diff --git a/rc_multiplayer/src/matches/aggregate.rs b/rc_multiplayer/src/matches/aggregate.rs new file mode 100644 index 0000000..4a59e63 --- /dev/null +++ b/rc_multiplayer/src/matches/aggregate.rs @@ -0,0 +1,72 @@ +pub struct GameMatches { + matches: std::collections::HashMap>, + routing: std::collections::HashMap, // user id to game guid +} + +impl GameMatches { + pub fn new() -> Self { + Self { + matches: std::collections::HashMap::new(), + routing: std::collections::HashMap::new(), + } + } + + pub fn spawn(self) -> tokio::sync::mpsc::Sender { + let (tx, rx) = tokio::sync::mpsc::channel(super::CHANNEL_BOUND); + tokio::spawn(self.run(rx)); + tx + } + + async fn start_new_match_engine(&self, _user: &Box, guid: &str) -> tokio::sync::mpsc::Sender { + // TODO figure out gamemode and act accordingly + let engine = super::GenericGamemodeEngine::new(guid.to_owned()); + engine.spawn() + } + + async fn run(mut self, mut rx: tokio::sync::mpsc::Receiver) { + log::info!("Match message router has started"); + while !rx.is_closed() { + if let Some(msg) = rx.recv().await { + log::debug!("Match message router got a message"); + match msg { + super::GameMessage::NewConnection { user, game_guid, connection, response, sender } => { + if let Some(tx) = self.matches.get(&game_guid) { + if tx.send(super::GameMessage::NewConnection { user, game_guid, connection, response, sender }).await.is_err() { + log::error!("Failed to send NewConnection game message to existing match"); + } + } else { + // create a new match + log::info!("Creating new game {}", game_guid); + let tx = self.start_new_match_engine(&user, &game_guid).await; + self.matches.insert(game_guid.clone(), tx.clone()); + self.routing.insert(user.user_id(), game_guid.clone()); + if tx.send(super::GameMessage::NewConnection { user, game_guid, connection, response, sender }).await.is_err() { + log::error!("Failed to send NewConnection game message to new match"); + } + } + } + msg => { + let user_id = msg.user_id(); + if let Some(guid) = self.routing.get(&user_id) { + if let Some(tx) = self.matches.get(guid) { + if tx.is_closed() { + self.matches.remove(guid); + self.routing.remove(&user_id); + } else { + if tx.send(msg).await.is_err() { + log::error!("Failed to route game message from user {} to match {}", user_id, guid); + } + } + } else { + self.routing.remove(&user_id); + } + } else { + log::warn!("Got unroutable user {}", user_id); + } + } + } + } + } + log::warn!("Match message router has completed"); + } +} diff --git a/rc_multiplayer/src/matches/engine.rs b/rc_multiplayer/src/matches/engine.rs new file mode 100644 index 0000000..ae2f7f7 --- /dev/null +++ b/rc_multiplayer/src/matches/engine.rs @@ -0,0 +1,3 @@ +pub trait GamemodeEngine: Send + Sync { + fn is_complete(&self) -> bool; +} diff --git a/rc_multiplayer/src/matches/generic.rs b/rc_multiplayer/src/matches/generic.rs new file mode 100644 index 0000000..5427d56 --- /dev/null +++ b/rc_multiplayer/src/matches/generic.rs @@ -0,0 +1,402 @@ +pub(super) struct UserConnection { + pub(super) user: std::sync::Arc>, + pub(super) connection: std::sync::Arc>, + pub(super) sender: std::sync::Arc>, + pub(super) state: UserState, + pub(super) machine: MachineState, +} + +pub(super) struct UserState { + pub(super) mode: std::sync::atomic::AtomicU8, + pub(super) progress: std::sync::atomic::AtomicU8, // percent + _x: (), +} + +impl UserState { + fn new() -> Self { + Self { + mode: std::sync::atomic::AtomicU8::new(ConnectionMode::Loading.to_u8()), + progress: std::sync::atomic::AtomicU8::new(0), + _x: (), + } + } +} + +pub(super) struct MachineState { + pub(super) selected_weapon: WeaponInfo, + _x: (), +} + +impl MachineState { + fn new() -> Self { + Self { + selected_weapon: WeaponInfo::new(), + _x: (), + } + } +} + +pub(super) struct WeaponInfo { + category: std::sync::atomic::AtomicU32, + size: std::sync::atomic::AtomicU32, +} + +impl WeaponInfo { + fn new() -> Self { + Self { + category: std::sync::atomic::AtomicU32::new(0), + size: std::sync::atomic::AtomicU32::new(0), + } + } +} + +#[repr(u8)] +#[derive(Debug, Copy, Clone)] +pub(super) enum ConnectionMode { + Loading = 0, + Sync = 1, + InGame = 2, +} + +impl ConnectionMode { + #[inline] + fn from_u8(num: u8) -> Self { + match num { + 0 => Self::Loading, + 1 => Self::Sync, + 2 => Self::InGame, + x => panic!("Unrecognized ConnectionMode {}", x), + } + } + + #[inline] + fn to_u8(self) -> u8 { + self as u8 + } +} + +pub(super) struct GenericGamemodeEngine { + pub users: tokio::sync::RwLock>, + pub user_id_map: tokio::sync::RwLock>, + //pub recv: tokio::sync::Mutex>, + //pub send: tokio::sync::mpsc::Sender, + pub game_guid: String, + pub is_complete: std::sync::atomic::AtomicBool, +} + +impl GenericGamemodeEngine { + pub fn new(guid: String) -> Self { + + Self { + users: tokio::sync::RwLock::new(std::collections::HashMap::new()), + user_id_map: tokio::sync::RwLock::new(std::collections::HashMap::new()), + game_guid: guid, + is_complete: std::sync::atomic::AtomicBool::new(false), + } + } + + pub(super) async fn user_key_by_user_id(&self, user_id: i32) -> Option { + self.user_id_map.read().await.get(&user_id).map(|x| *x) + } + + pub(super) async fn broadcast(&self, user_id: i32, code: rlnl::event_code::NetworkEvent, property: literustlib::packet::Property, data: T) { + for conn in self.users.read().await.values() { + if user_id == conn.user.user_id() { continue; } + let sender = crate::handlers::simple_typed::RlnlSender::new(&conn.sender); + crate::events::log_lnl_send_failure(sender.send_data( + &data, + code, + property, + &conn.connection, + ).await); + } + } + + pub(super) fn spawn(self) -> tokio::sync::mpsc::Sender { + let (tx, rx) = tokio::sync::mpsc::channel(super::CHANNEL_BOUND); + tokio::spawn(self.run(rx)); + tx + } + + pub(super) async fn run(self, mut recv: tokio::sync::mpsc::Receiver) { + while !recv.is_closed() { + if let Some(msg) = recv.recv().await { + match msg { + super::GameMessage::NewConnection { user, game_guid, connection, response, sender } => { + if self.game_guid != game_guid { + response.send(Some(super::messages::ErrorMessage { + message: "Game guid does not match".to_owned(), + inner: None, + })).unwrap_or_default(); + return; + } else { + let mut users = self.users.write().await; + let new_user = UserConnection { + user, + connection, + sender, + state: UserState::new(), + machine: MachineState::new(), + }; + //tokio::time::sleep(std::time::Duration::from_secs(1)).await; + let id = users.len() as u8; + if let Err(e) = self.send_loading_events(&new_user, id).await { + response.send(Some(super::messages::ErrorMessage { + message: "Failed to send GameGuidValidated response".to_owned(), + inner: Some(Box::new(e)), + })).unwrap_or_default(); + return; + } + log::debug!("User {} is validated to play game {}", new_user.user.user_id(), game_guid); + self.user_id_map.write().await.insert(new_user.user.user_id(), id); + users.insert(id, new_user); + response.send(None).unwrap_or_default(); + } + }, + super::GameMessage::LoadingProgress { user_id, user_name, progress } => { + let progress_data = rlnl::events::loading::LoadingProgress { + user_name: rlnl::types::BinaryWriterString(user_name), + progress, + }; + for conn in self.users.read().await.values() { + if user_id == conn.user.user_id() { + let progress_percent = (progress * 100.0).ceil() as u8; + log::debug!("User {} is loaded {}% into game {}", user_id, progress_percent, self.game_guid); + conn.state.progress.store(progress_percent, std::sync::atomic::Ordering::Relaxed); + } + let mode = ConnectionMode::from_u8(conn.state.mode.load(std::sync::atomic::Ordering::Relaxed)); + match mode { + ConnectionMode::Loading + | ConnectionMode::Sync => { + if user_id != conn.user.user_id() { + crate::events::log_lnl_send_failure(crate::handlers::simple_typed::RlnlSender::new(&conn.sender) + .send_data(&progress_data, rlnl::event_code::NetworkEvent::BroadcastLoadingProgress, literustlib::packet::Property::ReliableOrdered, &conn.connection).await); + } + /*if progress > 0.95 { + log::info!("User {} is ready, ending sync", user_id); + crate::events::log_lnl_send_failure(crate::handlers::simple_typed::RlnlSender::new(&conn.sender) + .send_empty( + rlnl::event_code::NetworkEvent::EndOfSync, + literustlib::packet::Property::ReliableOrdered, + &conn.connection, + ) + .await); + conn.mode.store(ConnectionMode::InGame.to_u8(), std::sync::atomic::Ordering::Relaxed); + }*/ + }, + ConnectionMode::InGame => { + log::warn!("Got loading progress for user {} who is supposed to be already in-game", user_id); + }, + } + } + } + super::GameMessage::RequestLoadingProgress { user_id } => { + let mut user_info = None; + for conn in self.users.read().await.values() { + if user_id == conn.user.user_id() { + user_info = Some(( + conn.sender.to_owned(), + conn.connection.to_owned(), + rlnl::events::loading::LoadingProgress { + user_name: rlnl::types::BinaryWriterString(conn.user.user_name().to_owned()), + progress: (conn.state.progress.load(std::sync::atomic::Ordering::Relaxed) as f32) / 100.0, + }, + )); + } + } + if let Some(user_info) = user_info { + let sender = crate::handlers::simple_typed::RlnlSender::new(&user_info.0); + for conn in self.users.read().await.values() { + if user_id == conn.user.user_id() { continue; } + crate::events::log_lnl_send_failure(sender.send_data( + &user_info.2, + rlnl::event_code::NetworkEvent::BroadcastLoadingProgress, + literustlib::packet::Property::ReliableOrdered, + &user_info.1, + ).await); + } + } else { + log::error!("Failed to find user {} in connected users for match {}", user_id, self.game_guid); + } + + }, + super::GameMessage::WeaponSelect { user_id, machine_id, category, size } => { + if let Some(conn) = self.users.read().await.get(&machine_id) { + let category_u32 = category as u32; + let size_u32 = size as u32; + conn.machine.selected_weapon.category.store(category_u32, std::sync::atomic::Ordering::Relaxed); + conn.machine.selected_weapon.size.store(size_u32, std::sync::atomic::Ordering::Relaxed); + let data = rlnl::events::ingame::SelectWeapon { + machine_id, + item_category: category_u32, + item_size: size_u32, + }; + self.broadcast( + user_id, + rlnl::event_code::NetworkEvent::BroadcastWeaponSelect, + literustlib::packet::Property::ReliableOrdered, + data, + ).await; + } + }, + super::GameMessage::RequestLoadingSync { user_id } => { + if let Some(user_key) = self.user_key_by_user_id(user_id).await { + if let Some(conn) = self.users.read().await.get(&user_key) { + self.spawn_send_sync_events(conn, user_id); + } + } + }, + super::GameMessage::Motion { user_id, data } => { + for conn in self.users.read().await.values() { + if conn.user.user_id() == user_id { continue; } // fun fact: the game hard crashes if you omit this + crate::events::log_lnl_send_failure(conn.sender.send_data(crate::handler::EventData { + message_ty: crate::data::MessageType::RobotMotion, + variant: 0, + data_size: data.len() as _, + data: data.clone(), + }, literustlib::packet::Property::Unreliable, &conn.connection).await); + } + } + super::GameMessage::NoOp => {}, + } + } + } + self.is_complete.store(true, std::sync::atomic::Ordering::Relaxed); + } + + async fn send_loading_events(&self, user: &UserConnection, player_id: u8) -> std::io::Result<()> { + let sender = crate::handlers::simple_typed::RlnlSender::new(&user.sender); + sender.send_data( + &rlnl::events::loading::PlayerID { owner: player_id }, + rlnl::event_code::NetworkEvent::GameGuidValidated, + literustlib::packet::Property::ReliableOrdered, + &user.connection + ).await?; + sender.send_data( + &rlnl::events::loading::PlayerIDsAndNames { + num_players: 2, + players: vec![ // FIXME + rlnl::events::loading::PlayerIDAndName { + player_id: 0, + name: rlnl::types::BinaryWriterString("NGniusness".to_owned()), + display_name: rlnl::types::BinaryWriterString("NGniusness".to_owned()), + }, + rlnl::events::loading::PlayerIDAndName { + player_id: 1, + name: rlnl::types::BinaryWriterString("NGniusness_echo".to_owned()), + display_name: rlnl::types::BinaryWriterString("NGniusness_echo".to_owned()), + }, + ], + }, + rlnl::event_code::NetworkEvent::PlayerIDs, + literustlib::packet::Property::ReliableOrdered, + &user.connection + ).await?; + sender.send_data( + &rlnl::events::loading::PlayerIDs { + num_ids: 0, + players: vec![], + }, + rlnl::event_code::NetworkEvent::HostAIs, + literustlib::packet::Property::ReliableOrdered, + &user.connection + ).await?; + Ok(()) + } + + fn spawn_send_sync_events(&self, user: &UserConnection, user_id: i32) { + let sender = user.sender.clone(); + let connection = user.connection.clone(); + tokio::spawn(Self::send_sync_events_wrapper(connection, sender, user_id)); + user.state.mode.store(ConnectionMode::Sync.to_u8(), std::sync::atomic::Ordering::Relaxed); + } + + async fn send_sync_events_wrapper(connection: std::sync::Arc>, sender: std::sync::Arc>, user_id: i32) { + if let Err(e) = Self::send_sync_events(connection, sender).await { + log::error!("Failed to send Sync events for user {}: {}", user_id, e); + } + } + + async fn send_sync_events(connection: std::sync::Arc>, sender: std::sync::Arc>) -> std::io::Result<()> { + let sender = crate::handlers::simple_typed::RlnlSender::new(&sender); + sender.send_empty( + rlnl::event_code::NetworkEvent::BeginSync, + literustlib::packet::Property::ReliableOrdered, + &connection) + .await?; + // sudden death + sender.send_data( + &rlnl::events::sync::UpdateGameModeSettings { // FIXME use value from config + respawn_heal_duration: 10.0, + respawn_full_heal_duration: 10.0, + }, + rlnl::event_code::NetworkEvent::GameModeSettings, + literustlib::packet::Property::ReliableOrdered, + &connection) + .await?; + sender.send_data( + &rlnl::events::GameTime(300.0), // FIXME use value from config + rlnl::event_code::NetworkEvent::CurrentGameTime, + literustlib::packet::Property::ReliableOrdered, + &connection) + .await?; + // generic + sender.send_data( + &rlnl::events::sync::InitialiseGameStats { + num_players: 2, + stats: vec![ // FIXME generate one per connection + rlnl::types::IngamePlayerStats { + player_name: 0, + num_stats: 0, + stats: vec![], + }, + rlnl::types::IngamePlayerStats { + player_name: 1, + num_stats: 0, + stats: vec![], + }, + ], + }, + rlnl::event_code::NetworkEvent::InitialiseGameStats, + literustlib::packet::Property::ReliableOrdered, + &connection) + .await?; + sender.send_data( + &rlnl::events::sync::SpawnPoint { + pos: rlnl::types::PosQuatPair { + pos: rlnl::types::CompressedVec3 { x: 0, y: 0, z: 0 }, + rot: rlnl::types::CompressedQuat { x: 0, y: 0, z: 0 }, + }, + owner: 0, + }, + rlnl::event_code::NetworkEvent::FreeSpawnPoint, + literustlib::packet::Property::ReliableOrdered, + &connection) + .await?; + /*sender.send_data( + &rlnl::events::sync::SyncMachineCubes { + machine_id: 0, + num_cubes: 0, + events: vec![ + rlnl::types::CubeState { + loc: rlnl::types::Byte3 { x: 0, y: 0, z: 0 }, + status: rlnl::types::CubeStatus { + ty: rlnl::types::CubeHistoryEventType::Heal, + damage: Some(1), + } + } + ], + }, + rlnl::event_code::NetworkEvent::SyncMachineCubes, + literustlib::packet::Property::ReliableOrdered, + &user.connection) + .await?;*/ + Ok(()) + } +} + +impl super::GamemodeEngine for GenericGamemodeEngine { + fn is_complete(&self) -> bool { + self.is_complete.load(std::sync::atomic::Ordering::Relaxed) // for now, this is never closed + } +} diff --git a/rc_multiplayer/src/matches/messages.rs b/rc_multiplayer/src/matches/messages.rs new file mode 100644 index 0000000..2dc30b9 --- /dev/null +++ b/rc_multiplayer/src/matches/messages.rs @@ -0,0 +1,65 @@ +pub enum GameMessage { + NewConnection { + user: std::sync::Arc>, + game_guid: String, + connection: std::sync::Arc>, + response: tokio::sync::oneshot::Sender>, + sender: std::sync::Arc>, + }, + LoadingProgress { + user_id: i32, + user_name: String, + progress: f32, + }, + RequestLoadingProgress { + user_id: i32, + }, + WeaponSelect { + user_id: i32, + machine_id: u8, + category: oj_rc_core::data::weapon_list::ItemCategory, + size: oj_rc_core::data::cube_list::ItemTier, + }, + RequestLoadingSync { + user_id: i32, + }, + Motion { + user_id: i32, + data: bytes::Bytes, + }, + NoOp, +} + +impl GameMessage { + pub fn user_id(&self) -> i32 { + match self { + Self::NewConnection { user, .. } => { + user.user_id() + } + Self::LoadingProgress { user_id, .. } => *user_id, + Self::RequestLoadingProgress { user_id, .. } => *user_id, + Self::WeaponSelect { user_id, .. } => *user_id, + Self::RequestLoadingSync { user_id, .. } => *user_id, + Self::Motion { user_id, .. } => *user_id, + Self::NoOp => unreachable!("NoOp is irrelevant for user ID"), + } + } +} + +#[derive(Debug)] +pub struct ErrorMessage { + pub message: String, + pub inner: Option>, +} + +impl std::error::Error for ErrorMessage {} + +impl core::fmt::Display for ErrorMessage { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + if let Some(inner) = &self.inner { + write!(f, "game communication error: {}; {}", self.message, inner) + } else { + write!(f, "game communication error: {}", self.message) + } + } +} diff --git a/rc_multiplayer/src/matches/mod.rs b/rc_multiplayer/src/matches/mod.rs new file mode 100644 index 0000000..7e8f590 --- /dev/null +++ b/rc_multiplayer/src/matches/mod.rs @@ -0,0 +1,13 @@ +mod engine; +pub use engine::GamemodeEngine; + +mod messages; +pub use messages::GameMessage; + +mod generic; +pub(self) use generic::GenericGamemodeEngine; + +mod aggregate; +pub use aggregate::GameMatches; + +pub const CHANNEL_BOUND: usize = 16; diff --git a/rc_multiplayer/src/traits.rs b/rc_multiplayer/src/traits.rs index 5bee4b7..036a2d6 100644 --- a/rc_multiplayer/src/traits.rs +++ b/rc_multiplayer/src/traits.rs @@ -3,9 +3,14 @@ pub type PacketData = crate::handler::EventData; #[async_trait::async_trait] pub trait EventCodeHandler: Send + Sync { - async fn handle(&self, data: &bytes::Bytes, peer: &std::sync::Arc>, user: &UserData, sender: &literustlib_server::DataSender); + async fn handle(&self, data: &bytes::Bytes, peer: &std::sync::Arc>, user: &UserData, sender: &std::sync::Arc>); } pub trait EventCode: EventCodeHandler { const CODE: i16; } + +#[async_trait::async_trait] +pub trait RobotMotionHandler: Send + Sync { + async fn handle(&self, data: &bytes::Bytes, user: &UserData); +} diff --git a/rc_multiplayer/src/user.rs b/rc_multiplayer/src/user.rs index f6d8ae7..6a69b6b 100644 --- a/rc_multiplayer/src/user.rs +++ b/rc_multiplayer/src/user.rs @@ -3,27 +3,61 @@ pub struct User { } impl User { - pub fn new() -> Self { + pub fn new(provider: std::sync::Arc) -> Self { Self { - state: tokio::sync::RwLock::new(UserState::Connecting), + state: tokio::sync::RwLock::new(UserState::Unauthenticated(provider)), } } pub async fn authenticate(&self, info: rlnl::events::loading::GameGuidInfo) -> bool { - *self.state.write().await = UserState::Authenticated(UserInfo { - game: info.game_guid.0, - username: info.player_name.0, - }); - true + let init_state_clone = self.state.read().await.clone(); + match init_state_clone { + UserState::Unauthenticated(auth) => { + let result = >::multiplayer_authenticate::<'_, '_>(&auth, info.player_name.0.clone()).await; + match result { + Ok(user) => { + *self.state.write().await = UserState::Authenticated(UserInfo { + game: info.game_guid.0, + user: std::sync::Arc::new(user), + }); + true + }, + Err(e) => { + log::error!("Failed to authenticate {}: {}", info.player_name.0, e); + false + } + } + }, + UserState::Authenticated(_) => { + log::warn!("User already authenticated, ignoring"); + true + } + } + } + + pub async fn user(&self) -> Option>> { + match &*self.state.read().await { + UserState::Unauthenticated(_) => None, + UserState::Authenticated(user) => Some(user.user.clone()), + } + } + + pub async fn game_guid(&self) -> Option { + match &*self.state.read().await { + UserState::Unauthenticated(_) => None, + UserState::Authenticated(user) => Some(user.game.clone()), + } } } +#[derive(Clone)] enum UserState { - Connecting, + Unauthenticated(std::sync::Arc), Authenticated(UserInfo), } +#[derive(Clone)] pub struct UserInfo { pub game: String, - pub username: String, + pub user: std::sync::Arc>, } diff --git a/rc_multiplayer/src/vehicle_motion.rs b/rc_multiplayer/src/vehicle_motion.rs new file mode 100644 index 0000000..9b72aca --- /dev/null +++ b/rc_multiplayer/src/vehicle_motion.rs @@ -0,0 +1,29 @@ +pub struct VehicleMotionHandler { + msg_router: tokio::sync::mpsc::Sender, +} + +pub(super) fn handler(init_ctx: &crate::InitConfig) -> VehicleMotionHandler { + VehicleMotionHandler::new(init_ctx) +} + +impl VehicleMotionHandler { + fn new(init_ctx: &crate::InitConfig) -> Self { + Self { + msg_router: init_ctx.matches_chann.clone(), + } + } +} + +#[async_trait::async_trait] +impl crate::RobotMotionHandler for VehicleMotionHandler { + async fn handle(&self, data: &bytes::Bytes, user: &crate::UserData) { + if let Some(user_info) = user.user().await { + crate::events::log_channel_send_failure(self.msg_router.send(crate::matches::GameMessage::Motion { + user_id: user_info.user_id(), + data: data.to_owned(), + }).await); + } else { + log::error!("Failed to handle motion unknown user"); + } + } +} diff --git a/rc_services_room/src/operations/mod.rs b/rc_services_room/src/operations/mod.rs index 4023836..2187669 100644 --- a/rc_services_room/src/operations/mod.rs +++ b/rc_services_room/src/operations/mod.rs @@ -217,4 +217,5 @@ pub fn handler(init_ctx: &crate::InitConfig) -> OperationsHandler .add(garage_slot_set_customisations::garage_slot_customisation_provider()) .add(garage_slot_name::garage_slot_rename_provider()) .add(garage_slot_copy::garage_slot_copy_provider()) + .add(polariton_server::operations::Ack::<12, _>::default()) // TODO handle UpdatePlayerDailyQuestProgressRequest instead of ignoring it } diff --git a/rc_services_room/src/operations/singleplayer_campaigns.rs b/rc_services_room/src/operations/singleplayer_campaigns.rs index 33e2438..ed32f58 100644 --- a/rc_services_room/src/operations/singleplayer_campaigns.rs +++ b/rc_services_room/src/operations/singleplayer_campaigns.rs @@ -88,7 +88,7 @@ pub(super) fn singleplayer_save_complete_campaign_provider() -> SimpleFunc<68, c if let Some(Typed::Int(campaign_difficulty)) = params.get(&CAMPAIGN_DIFFICULTY_PARAM_KEY) { if let Some(Typed::Int(wave_number)) = params.get(&CAMPAIGN_WAVE_NUMBER_PARAM_KEYL) { let user_info = user.user()?; - log::info!("User {} completed campaign {} difficulty {} wave {}", user_info.token().uuid, campaign_id.string, campaign_difficulty, wave_number); + log::info!("User {} completed campaign {} difficulty {} wave {}", user_info.public_id(), campaign_id.string, campaign_difficulty, wave_number); // TODO save wave as completed params.clear(); }