diff --git a/Cargo.lock b/Cargo.lock index a9cd36c..dc3cb96 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3074,7 +3074,7 @@ checksum = "059c95245738cdc7b40078cdd51a23200252a4c0a0a6dd005136152b3f467a4a" [[package]] name = "oj_cdn" -version = "1.3.1" +version = "1.4.0" dependencies = [ "actix-files", "actix-web", @@ -3103,7 +3103,7 @@ dependencies = [ [[package]] name = "oj_polariton_auth" -version = "1.3.1" +version = "1.4.0" dependencies = [ "log", "num", @@ -3115,7 +3115,7 @@ dependencies = [ [[package]] name = "oj_rc_auth" -version = "1.3.1" +version = "1.4.0" dependencies = [ "actix-files", "actix-web", @@ -3140,7 +3140,7 @@ dependencies = [ [[package]] name = "oj_rc_chat" -version = "1.3.1" +version = "1.4.0" dependencies = [ "clap", "env_logger", @@ -3153,7 +3153,7 @@ dependencies = [ [[package]] name = "oj_rc_chat_room" -version = "1.3.1" +version = "1.4.0" dependencies = [ "async-trait", "chrono", @@ -3176,7 +3176,7 @@ dependencies = [ [[package]] name = "oj_rc_core" -version = "1.3.1" +version = "1.4.0" dependencies = [ "argon2", "async-trait", @@ -3208,7 +3208,7 @@ dependencies = [ [[package]] name = "oj_rc_database" -version = "1.3.1" +version = "1.4.0" dependencies = [ "async-trait", "itertools 0.14.0", @@ -3221,7 +3221,7 @@ dependencies = [ [[package]] name = "oj_rc_factory" -version = "1.3.1" +version = "1.4.0" dependencies = [ "async-trait", "base64 0.22.1", @@ -3236,7 +3236,7 @@ dependencies = [ [[package]] name = "oj_rc_factory_web" -version = "1.3.1" +version = "1.4.0" dependencies = [ "actix-files", "actix-web", @@ -3255,7 +3255,7 @@ dependencies = [ [[package]] name = "oj_rc_lobby" -version = "1.3.1" +version = "1.4.0" dependencies = [ "clap", "env_logger", @@ -3268,7 +3268,7 @@ dependencies = [ [[package]] name = "oj_rc_lobby_room" -version = "1.3.1" +version = "1.4.0" dependencies = [ "async-trait", "chrono", @@ -3289,7 +3289,7 @@ dependencies = [ [[package]] name = "oj_rc_microtransactions" -version = "1.3.1" +version = "1.4.0" dependencies = [ "actix-web", "clap", @@ -3303,7 +3303,7 @@ dependencies = [ [[package]] name = "oj_rc_multiplayer" -version = "1.3.1" +version = "1.4.0" dependencies = [ "async-trait", "atomic_float", @@ -3326,7 +3326,7 @@ dependencies = [ [[package]] name = "oj_rc_plugins" -version = "1.3.1" +version = "1.4.0" dependencies = [ "libloading", "log", @@ -3334,7 +3334,7 @@ dependencies = [ [[package]] name = "oj_rc_services" -version = "1.3.1" +version = "1.4.0" dependencies = [ "clap", "env_logger", @@ -3347,7 +3347,7 @@ dependencies = [ [[package]] name = "oj_rc_services_room" -version = "1.3.1" +version = "1.4.0" dependencies = [ "async-trait", "chrono", @@ -3373,7 +3373,7 @@ dependencies = [ [[package]] name = "oj_rc_singleplayer" -version = "1.3.1" +version = "1.4.0" dependencies = [ "clap", "env_logger", @@ -3386,7 +3386,7 @@ dependencies = [ [[package]] name = "oj_rc_singleplayer_room" -version = "1.3.1" +version = "1.4.0" dependencies = [ "async-trait", "chrono", @@ -3404,7 +3404,7 @@ dependencies = [ [[package]] name = "oj_rc_social" -version = "1.3.1" +version = "1.4.0" dependencies = [ "clap", "env_logger", @@ -3417,7 +3417,7 @@ dependencies = [ [[package]] name = "oj_rc_social_room" -version = "1.3.1" +version = "1.4.0" dependencies = [ "async-trait", "chrono", @@ -3435,7 +3435,7 @@ dependencies = [ [[package]] name = "oj_rc_society" -version = "1.3.1" +version = "1.4.0" dependencies = [ "actix-files", "actix-identity", diff --git a/rc_auth/src/robocraft/intercom/lobby.rs b/rc_auth/src/robocraft/intercom/lobby.rs index 66e985c..cc2bd09 100644 --- a/rc_auth/src/robocraft/intercom/lobby.rs +++ b/rc_auth/src/robocraft/intercom/lobby.rs @@ -5,12 +5,45 @@ use actix_web::{rt, web::{Payload, Data, Json, Path}, Error, HttpRequest, HttpRe #[get("/intercom/userless/.oj_lobby")] pub async fn lobby_state_ws(req: HttpRequest, stream: Payload, auth: Data, reg: Data) -> Result { auth.validate(&req, "state/.oj_lobby")?; - let (res, mut session, _stream) = actix_ws::handle(&req, stream)?; + let (res, mut session, stream) = actix_ws::handle(&req, stream)?; - /*let mut stream = stream + let mut stream = stream .aggregate_continuations() .max_continuation_size(2_usize.pow(20)); // aggregate continuation frames up to 1MiB - */ + + let mut session_rx = session.clone(); + // start receiver task but don't wait for it + rt::spawn(async move { + while let Some(msg_result) = stream.recv().await { + match msg_result { + Ok(msg) => { + match msg { + actix_ws::AggregatedMessage::Text(s) => { + // TODO allow responses + log::debug!("Received `{}` ({}) from lobby websocket receiver", s, s.len()); + }, + actix_ws::AggregatedMessage::Binary(_) => {}, + actix_ws::AggregatedMessage::Ping(p) => { + log::debug!("Received {}B ping from lobby websocket", p.len()); + session_rx.pong(&p).await.unwrap_or_default(); + }, + actix_ws::AggregatedMessage::Pong(_) => {}, + actix_ws::AggregatedMessage::Close(reason) => { + if let Some(reason) = reason { + log::info!("Closed lobby websocket receiver with reason {:?}", reason); + } else { + log::warn!("Closed lobby websocket receiver without reason"); + } + break; + }, + } + }, + Err(e) => { + log::error!("Received lobby websocket error: {}", e); + }, + } + } + }); let (tx, mut rx) = tokio::sync::mpsc::channel(16); reg.register_lobby_state_service(tx).await; @@ -19,6 +52,7 @@ pub async fn lobby_state_ws(req: HttpRequest, stream: Payload, auth: Data { diff --git a/rc_auth/src/robocraft/intercom/services.rs b/rc_auth/src/robocraft/intercom/services.rs index b3587d0..adf054f 100644 --- a/rc_auth/src/robocraft/intercom/services.rs +++ b/rc_auth/src/robocraft/intercom/services.rs @@ -5,20 +5,58 @@ use actix_web::{rt, web::{Payload, Data, Path, Json}, Error, HttpRequest, HttpRe #[get("/intercom/.oj_services/{name}")] pub async fn services_ws(req: HttpRequest, stream: Payload, auth: Data, reg: Data, name: Path) -> Result { auth.validate(&req, &format!(".oj_services/{}", urlencoding::encode(&name)))?; - let (res, mut session, _stream) = actix_ws::handle(&req, stream)?; + let (res, mut session, stream) = actix_ws::handle(&req, stream)?; - /*let mut stream = stream + let mut stream = stream .aggregate_continuations() .max_continuation_size(2_usize.pow(20)); // aggregate continuation frames up to 1MiB - */ let (tx, mut rx) = tokio::sync::mpsc::channel(16); - reg.register_web_service(name.clone(), tx).await; + reg.register_web_service(name.clone(), tx.clone()).await; log::debug!("Registered web services intercom websocket for user {}", name); - // start task but don't wait for it + let name_rx: String = name.clone(); + let mut session_rx = session.clone(); + // start receiver task but don't wait for it + rt::spawn(async move { + while let Some(msg_result) = stream.recv().await { + match msg_result { + Ok(msg) => { + match msg { + actix_ws::AggregatedMessage::Text(s) => { + // TODO allow responses + log::debug!("Received `{}` ({}) from services websocket receiver for user {}", s, s.len(), name_rx); + }, + actix_ws::AggregatedMessage::Binary(_) => {}, + actix_ws::AggregatedMessage::Ping(p) => { + log::debug!("Received {}B ping from services websocket for user {}: {}", p.len(), name_rx, String::from_utf8_lossy(&p)); + session_rx.pong(&p).await.unwrap_or_default(); + }, + actix_ws::AggregatedMessage::Pong(p) => { + log::debug!("Received {}B pong from services websocket for user {}: {}", p.len(), name_rx, String::from_utf8_lossy(&p)); + }, + actix_ws::AggregatedMessage::Close(reason) => { + if let Some(reason) = reason { + log::info!("Closed services websocket receiver for user {} with reason {:?}", name_rx, reason); + } else { + log::warn!("Closed services websocket receiver for user {} without reason", name_rx); + } + break; + }, + } + }, + Err(e) => { + log::warn!("Received services websocket error for user {}: {}", name_rx, e); + }, + } + } + tx.send(crate::robocraft::intercom::IntercomOp::Info(crate::robocraft::intercom::IntercomInfo::Close)).await.unwrap_or_default(); + }); + + // start sender task but don't wait for it rt::spawn(async move { let mut is_ok = false; + session.ping(b"openjam").await.unwrap_or_default(); while let Some(op) = rx.recv().await { match op { super::IntercomOp::Message(msg) => { diff --git a/rc_core/src/persist/user/intercom.rs b/rc_core/src/persist/user/intercom.rs index e649713..702ddc5 100644 --- a/rc_core/src/persist/user/intercom.rs +++ b/rc_core/src/persist/user/intercom.rs @@ -1,7 +1,121 @@ use serde::{Serialize, Deserialize}; +pub struct IntercomListener { + websocket: reqwest_websocket::WebSocket, + _d: std::marker::PhantomData, +} + +impl IntercomListener { + pub(super) fn new(websocket: reqwest_websocket::WebSocket) -> Self { + Self { + websocket, + _d: Default::default(), + } + } + + /*pub async fn close(self, reason: &str) { + if let Err(e) = self.websocket.close(reqwest_websocket::CloseCode::Normal, Some(reason)).await { + log::error!("Failed to close intercom websocket with reason {}: {}", reason, e); + } + }*/ +} + +impl IntercomListener { + pub async fn listen(self) -> impl futures::Stream> + Unpin { + use futures::StreamExt; + //self.websocket.map(|msg_res| msg_res.and_then(|msg| msg.json())) + let (sink, stream) = self.websocket.split(); + let sink = std::sync::Arc::new(tokio::sync::Mutex::new(sink)); + stream.filter_map(move |msg_res| { + let sink = sink.clone(); + Box::pin(async move { + let res = match msg_res.ok() { + None => None, + /*reqwest_websocket::Message::Text(s) => { + (serde_json::from_str(&s).map_err(reqwest_websocket::Error::Json), None) + },*/ + Some(reqwest_websocket::Message::Ping(p)) => { + use futures::SinkExt; + log::debug!("Received {}B intercom websocket ping: {}", p.len(), String::from_utf8_lossy(&p)); + let mut sink_lock = sink.lock().await; + if let Err(e) = sink_lock.send(reqwest_websocket::Message::Pong(p)).await { + log::warn!("Failed to send intercom websocket pong: {}", e); + } + None + }, + Some(reqwest_websocket::Message::Pong(p)) => { + log::debug!("Received {}B intercom websocket pong: {}", p.len(), String::from_utf8_lossy(&p)); + None + }, + Some(msg) => Some(msg.json()) + }; + res + } + )}) + } +} + +impl IntercomListener { + pub async fn send(self) -> impl futures::Sink { + use futures::SinkExt; + self.websocket.with::(|d| { + async move { + serde_json::to_string(&d) + .map_err(reqwest_websocket::Error::Json) + .map(reqwest_websocket::Message::Text) + } + }) + } +} + +impl IntercomListener { + pub async fn split(self) -> (impl futures::Sink, impl futures::Stream> + Unpin) { + use futures::SinkExt; + use futures::StreamExt; + let (sink, stream) = self.websocket.split(); + let sink = std::sync::Arc::new(tokio::sync::Mutex::new(sink)); + let sink_impl_instance = sink.clone(); + let sink_impl = futures::sink::unfold(sink_impl_instance, |sink_instance, d| async move { + let msg = serde_json::to_string(&d) + .map_err(reqwest_websocket::Error::Json) + .map(reqwest_websocket::Message::Text)?; + let mut lock = sink_instance.lock().await; + lock.send(msg).await?; + drop(lock); + Ok::<_, reqwest_websocket::Error>(sink_instance) + }); + let stream_impl = stream.filter_map(move |msg_res| { + let sink = sink.clone(); + Box::pin(async move { + let res = match msg_res.ok() { + None => None, + /*reqwest_websocket::Message::Text(s) => { + (serde_json::from_str(&s).map_err(reqwest_websocket::Error::Json), None) + },*/ + Some(reqwest_websocket::Message::Ping(p)) => { + use futures::SinkExt; + log::debug!("Received {}B intercom websocket ping: {}", p.len(), String::from_utf8_lossy(&p)); + let mut sink_lock = sink.lock().await; + if let Err(e) = sink_lock.send(reqwest_websocket::Message::Pong(p)).await { + log::warn!("Failed to send intercom websocket pong: {}", e); + } + None + }, + Some(reqwest_websocket::Message::Pong(p)) => { + log::debug!("Received {}B intercom websocket pong: {}", p.len(), String::from_utf8_lossy(&p)); + None + }, + Some(msg) => Some(msg.json()) + }; + res + } + )}); + (sink_impl, stream_impl) + } +} + impl super::account_json::UserData { - async fn listen_on_websocket(&self, server_name: &str) -> Result, reqwest_websocket::Error> { + async fn listen_on_websocket(&self, server_name: &str) -> Result, reqwest_websocket::Error> { use reqwest_websocket::Upgrade; let url_encoded_pub_id = urlencoding::encode(&self.account.public_id); let token = generate_token(format!("{}/{}", server_name, url_encoded_pub_id).as_bytes(), &self.secret); @@ -15,10 +129,7 @@ impl super::account_json::UserData { .await? .into_websocket() .await?; - Ok(super::IntercomListener { - websocket, - _d: Default::default(), - }) + Ok(IntercomListener::new(websocket)) } async fn post_to_intercom(&self, data: &D, server_name: &str, operation: &str) -> Result<(), reqwest::Error> { @@ -86,7 +197,7 @@ impl super::IntercomUser for super::account_json::UserData { Ok(()) } - async fn webservice_listener(&self) -> Result, polariton_server::operations::SimpleOpError> { + async fn webservice_listener(&self) -> Result, polariton_server::operations::SimpleOpError> { self.listen_on_websocket(".oj_services").await .map_err(|e| polariton_server::operations::SimpleOpError::with_message( crate::data::error_codes::WebServicesError::PlatformFeatureNotAvailable as i16, diff --git a/rc_core/src/persist/user/mod.rs b/rc_core/src/persist/user/mod.rs index 3baaebb..a611718 100644 --- a/rc_core/src/persist/user/mod.rs +++ b/rc_core/src/persist/user/mod.rs @@ -11,10 +11,11 @@ mod inventory; pub use inventory::{UnlockedParts, UnlockOverride}; mod traits; -pub use traits::{UserProvider, User, UserToken, UserSlots, UserSlotData, VehicleData, UserAuthInfo, FederatedAuthInfo, UserLoginInfo, UserAuthenticator, FederatedAuthenticator, NewSlotData, UserId, RegistrationInfo, VehicleUploadData, ChatUser, AvatarInfo, GetAvatarInfo, ControlData, ControlType, CustomisationData, GetCustomisationData, SetSanction, SanctionType, LobbyUser, GameDescriptor, PlayerLobbyDescriptor, MultiplayerUser, PlayerScore, MultiplayerError, MultiplayerErrorCode, PlayerDescriptor, GameEventSetter, CurrentGameEvent, AuthError, IntercomUser, FakePlayers, ResolvedVehicle, CommonUser, IntercomListener, UserRole, SocialUser, SocialUserC, CurrencyType, CurrencyOp, MatchRewards, SingleplayerUser, PurchaseResult, FactoryUser, FriendInviteReturn, FriendData, FriendInviteStatus, SocialInfo, ClanData, ClanMember, ClanMemberRank, ClanType, ClanSearchQuery, ClanInviteData, Userless, GameOverrides, WebUser, GarageWebInfo, GarageWebStats, SanctionWebStats, AccountWebStats, SocialWebStats, FederationWebData, FederationWebDetails}; +pub use traits::{UserProvider, User, UserToken, UserSlots, UserSlotData, VehicleData, UserAuthInfo, FederatedAuthInfo, UserLoginInfo, UserAuthenticator, FederatedAuthenticator, NewSlotData, UserId, RegistrationInfo, VehicleUploadData, ChatUser, AvatarInfo, GetAvatarInfo, ControlData, ControlType, CustomisationData, GetCustomisationData, SetSanction, SanctionType, LobbyUser, GameDescriptor, PlayerLobbyDescriptor, MultiplayerUser, PlayerScore, MultiplayerError, MultiplayerErrorCode, PlayerDescriptor, GameEventSetter, CurrentGameEvent, AuthError, IntercomUser, FakePlayers, ResolvedVehicle, CommonUser, UserRole, SocialUser, SocialUserC, CurrencyType, CurrencyOp, MatchRewards, SingleplayerUser, PurchaseResult, FactoryUser, FriendInviteReturn, FriendData, FriendInviteStatus, SocialInfo, ClanData, ClanMember, ClanMemberRank, ClanType, ClanSearchQuery, ClanInviteData, Userless, GameOverrides, WebUser, GarageWebInfo, GarageWebStats, SanctionWebStats, AccountWebStats, SocialWebStats, FederationWebData, FederationWebDetails}; pub mod intercom; pub use intercom::generate_token as generate_intercom_token; +pub use intercom::IntercomListener; mod multiplayer; mod lobby; diff --git a/rc_core/src/persist/user/traits.rs b/rc_core/src/persist/user/traits.rs index f1bb124..3de5412 100644 --- a/rc_core/src/persist/user/traits.rs +++ b/rc_core/src/persist/user/traits.rs @@ -431,7 +431,7 @@ pub struct PlayerScore { pub trait IntercomUser: CommonUser { async fn save_custom_avatar(&self, image: Vec) -> Result<(), polariton_server::operations::SimpleOpError>; async fn save_factory_thumbnail(&self, factory_id: i32, image: Vec) -> Result<(), polariton_server::operations::SimpleOpError>; - async fn webservice_listener(&self) -> Result, polariton_server::operations::SimpleOpError>; + async fn webservice_listener(&self) -> Result, polariton_server::operations::SimpleOpError>; async fn show_dev_message(&self, msg: super::intercom::IntercomDevMessage, to: Vec); async fn enter_maintenance(&self, msg: super::intercom::IntercomMaintenanceMessage, to: Vec); async fn trigger_workaround(&self, msg: super::intercom::IntercomWorkaroundMessage, to: Vec); @@ -439,21 +439,6 @@ pub trait IntercomUser: CommonUser { async fn update_status(&self, server_name: &str, msg: oj_serdes::ServerStatus); } -pub struct IntercomListener { - pub(super) websocket: reqwest_websocket::WebSocket, - pub(super) _d: std::marker::PhantomData, -} - -impl IntercomListener { - pub async fn listen(self) -> impl futures::Stream>> + Unpin { - use futures::StreamExt; - self.websocket.map(|msg| - msg.map_err(Box::new) - .and_then(|msg| msg.json().map_err(Box::new)) - ) - } -} - pub struct ResolvedVehicle { pub mastery: i32, pub tier: i32, diff --git a/rc_core/src/persist/user/userless.rs b/rc_core/src/persist/user/userless.rs index ea98981..a70d655 100644 --- a/rc_core/src/persist/user/userless.rs +++ b/rc_core/src/persist/user/userless.rs @@ -14,10 +14,7 @@ impl super::AccountProvider { .await? .into_websocket() .await?; - Ok(super::IntercomListener { - websocket, - _d: Default::default(), - }) + Ok(super::IntercomListener::new(websocket)) } }