use std::net::SocketAddr; use std::sync::Arc; use std::time::Duration; use tokio::task::JoinHandle; use logic::{ message_type, CancelMatchmakeDoneMessage, ClientCapabilitiesMessage, GoHomeMessage, KeepAliveServerMessage, LoginFailedMessage, LoginMessage, LoginOkMessage, ServerErrorMessage, }; use service_rpc::{ AccountRef, DeviceInfo, HomeRequestKind, LoginOutcome, Session as AuthSession, WireMessage, }; use titan::{Incoming, LogicLong, MessagingSender, OutboundMessage}; use crate::backend::Backends; use crate::config::GatewayConfig; #[derive(Debug, thiserror::Error)] pub enum RoutingError { #[error("transport: {0}")] Transport(#[from] titan::MessagingError), #[error("backend: {0}")] Backend(#[from] service_rpc::RpcError), #[error("client sent message type {0} before authenticating")] Unauthenticated(u16), } type RoutingResult = Result; pub const BATTLE_TICK_INTERVAL: Duration = Duration::from_millis(50); pub struct MessageManager { peer: SocketAddr, config: Arc, backends: Arc, sender: MessagingSender, account: Option, battle_ticker: Option>, } impl Drop for MessageManager { fn drop(&mut self) { self.stop_battle_ticker(); } } impl MessageManager { pub fn new( peer: SocketAddr, config: Arc, backends: Arc, sender: MessagingSender, ) -> Self { Self { peer, config, backends, sender, account: None, battle_ticker: None, } } pub fn account(&self) -> Option { self.account } pub async fn receive_message(&mut self, incoming: Incoming) -> RoutingResult<()> { let message_type = incoming.message_type(); if self.account.is_none() && message_type != message_type::LOGIN { return Err(RoutingError::Unauthenticated(message_type)); } let Some(message) = incoming.message.as_ref() else { return Ok(()); }; if !message.direction().accepts_from_client() { tracing::warn!( peer = %self.peer, message = message.name(), "client sent a server to client message" ); return Ok(()); } match message_type { message_type::LOGIN => { let login = incoming.downcast::().expect("login message"); self.handle_login(login).await } message_type::KEEP_ALIVE => { self.sender .send_message(&KeepAliveServerMessage::default()) .await?; Ok(()) } message_type::CLIENT_CAPABILITIES => { let capabilities = incoming .downcast::() .expect("client capabilities message"); if let Some(account) = self.account { self.backends .game .client_capabilities(account, capabilities.ping_ms) .await?; } Ok(()) } message_type::GO_HOME => { let _ = incoming.downcast::(); self.stop_battle_ticker(); self.push_home(HomeRequestKind::GoHome).await } message_type::SECTOR_COMMAND => { let Some(account) = self.account else { return Ok(()); }; if let Err(error) = self .backends .game .sector_command(account, incoming.payload) .await { tracing::warn!(peer = %self.peer, %error, "could not deliver the sector command"); } Ok(()) } message_type::SEND_BATTLE_EVENT => { let Some(account) = self.account else { return Ok(()); }; if let Err(error) = self .backends .game .battle_event(account, incoming.payload) .await { tracing::warn!(peer = %self.peer, %error, "could not deliver the battle event"); } Ok(()) } message_type::CANCEL_MATCHMAKE => { tracing::info!(peer = %self.peer, "client cancelled matchmaking"); self.stop_battle_ticker(); if let Some(account) = self.account { let _ = self.backends.game.cancel_matchmake(account).await; } self.sender .send_message(&CancelMatchmakeDoneMessage::default()) .await?; Ok(()) } message_type::HOME_LOGIC_STOPPED => { let Some(account) = self.account else { return Ok(()); }; match self .backends .game .home_logic_stopped(account, incoming.payload) .await { Ok(replies) => { self.push_wire_messages(replies).await?; self.start_battle_ticker(); } Err(error) => { tracing::warn!(peer = %self.peer, %error, "could not match the player"); self.sender .send_message(&CancelMatchmakeDoneMessage::default()) .await?; } } Ok(()) } message_type::START_MISSION => { let Some(account) = self.account else { return Ok(()); }; match self .backends .game .start_mission(account, incoming.payload) .await { Ok(replies) => { self.push_wire_messages(replies).await?; self.start_battle_ticker(); } Err(error) => { tracing::warn!(peer = %self.peer, %error, "could not start the mission"); self.sender .send_message(&ServerErrorMessage::new( "this server could not build the battle", )) .await?; } } Ok(()) } message_type::END_CLIENT_TURN => { if let Some(account) = self.account { let replies = self .backends .game .end_client_turn(account, incoming.payload) .await?; self.push_wire_messages(replies).await?; } Ok(()) } _ => { tracing::debug!( peer = %self.peer, message = message.name(), "no gateway route for this message" ); Ok(()) } } } pub async fn on_disconnect(&self) { if let Some(account) = self.account { if let Err(error) = self.backends.game.disconnect(account).await { tracing::debug!(peer = %self.peer, %error, "disconnect notification failed"); } } } async fn handle_login(&mut self, login: &LoginMessage) -> RoutingResult<()> { let account = AccountRef::new(login.account_id.high, login.account_id.low); let device = DeviceInfo { udid: login.udid.clone(), open_udid: login.open_udid.clone(), device: login.device.clone(), os_version: login.os_version.clone(), android: login.android, preferred_language: login.preferred_device_language.clone(), client_major_version: login.client_major_version, client_minor_version: login.client_minor_version, client_build: login.client_build, resource_sha: login.resource_sha.clone(), }; tracing::info!( peer = %self.peer, %account, version = format!( "{}.{}.{}", login.client_major_version, login.client_minor_version, login.client_build ), device = login.device.as_deref().unwrap_or("unknown"), "login attempt" ); match self .backends .auth .login(account, login.pass_token.clone(), device) .await? { LoginOutcome::Rejected { error_code, message, } => { tracing::info!( peer = %self.peer, error_code, reason = message.as_deref().unwrap_or("unspecified"), "login rejected" ); self.sender .send_message(&LoginFailedMessage { error_code, fingerprint_json: None, redirect_domain: None, content_url: None, remaining_seconds: 0, }) .await?; Ok(()) } LoginOutcome::Accepted(session) => { self.account = Some(session.account); self.sender.send_message(&self.login_ok(&session)).await?; self.push_home(HomeRequestKind::Login).await } } } fn login_ok(&self, session: &AuthSession) -> LoginOkMessage { let account = LogicLong::new(session.account.high, session.account.low); LoginOkMessage { account_id: account, home_id: account, pass_token: Some(session.pass_token.clone()), facebook_id: None, gamecenter_id: None, server_major_version: self.config.server_major_version, server_minor_version: self.config.server_minor_version, server_build: self.config.server_build, content_version: self.config.content_version, server_environment: Some(self.config.server_environment.clone()), session_count: session.session_count, play_time_seconds: session.play_time_seconds, days_since_started_playing: session.days_since_started_playing, facebook_app_id: self.config.facebook_app_id.clone(), server_time: Some(session.server_time.clone()), account_created_date: Some(session.account_created_date.clone()), startup_cooldown_seconds: self.config.startup_cooldown_seconds, } } async fn push_home(&self, kind: HomeRequestKind) -> RoutingResult<()> { let Some(account) = self.account else { return Ok(()); }; let messages = self.backends.game.load_home(account, kind).await?; self.push_wire_messages(messages).await } fn start_battle_ticker(&mut self) { let Some(account) = self.account else { return; }; self.stop_battle_ticker(); let backends = Arc::clone(&self.backends); let sender = self.sender.clone(); let peer = self.peer; self.battle_ticker = Some(tokio::spawn(async move { let mut ticker = tokio::time::interval(BATTLE_TICK_INTERVAL); ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); loop { ticker.tick().await; if sender.is_closed() { return; } match backends.game.battle_tick(account).await { Ok(messages) => { for message in messages { let outbound = OutboundMessage::raw( message.message_type, message.message_version, message.payload, ); if sender.send(outbound).await.is_err() { return; } } } Err(error) => { tracing::warn!(%peer, %error, "the battle tick failed, stopping the ticker"); return; } } } })); } fn stop_battle_ticker(&mut self) { if let Some(handle) = self.battle_ticker.take() { handle.abort(); } } async fn push_wire_messages(&self, messages: Vec) -> RoutingResult<()> { for message in messages { self.sender .send(OutboundMessage::raw( message.message_type, message.message_version, message.payload, )) .await?; } Ok(()) } }