diff --git a/Cargo.lock b/Cargo.lock index b156014..cd77382 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2558,6 +2558,7 @@ dependencies = [ "libfj", "log", "oj_rc_core", + "oj_serdes", "serde", "serde_json", "tokio", @@ -2588,6 +2589,7 @@ dependencies = [ "log", "oj_polariton_auth", "oj_rc_core", + "oj_serdes", "polariton", "polariton_server", "regex", @@ -2612,6 +2614,7 @@ dependencies = [ "num-quaternion", "oj_rc_database", "oj_rc_factory", + "oj_serdes", "polariton", "polariton_server", "rand 0.9.0", @@ -2670,6 +2673,7 @@ dependencies = [ "oj_polariton_auth", "oj_rc_core", "oj_rc_factory", + "oj_serdes", "polariton", "polariton_server", "tokio", @@ -2702,6 +2706,7 @@ dependencies = [ "log", "num-quaternion", "oj_rc_core", + "oj_serdes", "rand 0.9.0", "rlnl", "tokio", @@ -2735,6 +2740,7 @@ dependencies = [ "oj_polariton_auth", "oj_rc_core", "oj_rc_factory", + "oj_serdes", "polariton", "polariton_server", "rand 0.9.0", @@ -2761,11 +2767,13 @@ name = "oj_rc_singleplayer_room" version = "0.5.0" dependencies = [ "async-trait", + "chrono", "clap", "env_logger", "log", "oj_polariton_auth", "oj_rc_core", + "oj_serdes", "polariton", "polariton_server", "tokio", @@ -2789,16 +2797,26 @@ name = "oj_rc_social_room" version = "0.5.0" dependencies = [ "async-trait", + "chrono", "clap", "env_logger", "log", "oj_polariton_auth", "oj_rc_core", + "oj_serdes", "polariton", "polariton_server", "tokio", ] +[[package]] +name = "oj_serdes" +version = "0.1.0" +dependencies = [ + "serde", + "serde_json", +] + [[package]] name = "once_cell" version = "1.20.2" diff --git a/Cargo.toml b/Cargo.toml index ad0c16c..dbc62a0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -49,3 +49,4 @@ rand = { version = "0.9", features = [ "thread_rng" ] } num-quaternion = "1.0" hex = "0.4" base64 = "0.22" +oj_serdes = { version = "0.1.0", path = "../oj_core/serdes" } diff --git a/rc_auth/Cargo.toml b/rc_auth/Cargo.toml index f825f3f..cc241f3 100644 --- a/rc_auth/Cargo.toml +++ b/rc_auth/Cargo.toml @@ -22,5 +22,6 @@ git-version.workspace = true serde.workspace = true serde_json.workspace = true hex.workspace = true +oj_serdes.workspace = true handlebars = { version = "5", features = ["dir_source"] } diff --git a/rc_auth/src/main.rs b/rc_auth/src/main.rs index b4f6f9d..d8cf5ea 100644 --- a/rc_auth/src/main.rs +++ b/rc_auth/src/main.rs @@ -64,6 +64,8 @@ async fn main() -> std::io::Result<()> { .service(robocraft::username::user_password_auth) .service(robocraft::intercom::services_ws) .service(robocraft::intercom::service_msg) + .service(robocraft::intercom::status_get) + .service(robocraft::intercom::status_set) }) .bind((cli_args.ip, cli_args.port))? .run() diff --git a/rc_auth/src/robocraft/intercom/mod.rs b/rc_auth/src/robocraft/intercom/mod.rs index d95e68f..302525b 100644 --- a/rc_auth/src/robocraft/intercom/mod.rs +++ b/rc_auth/src/robocraft/intercom/mod.rs @@ -7,6 +7,9 @@ pub use services::{services_ws, service_msg}; mod user_registry; pub use user_registry::Users; +mod status; +pub use status::{status_set, status_get}; + enum IntercomOp { Message(oj_rc_core::persist::user::intercom::IntercomWebServiceUserMessage), Info(IntercomInfo), diff --git a/rc_auth/src/robocraft/intercom/status.rs b/rc_auth/src/robocraft/intercom/status.rs new file mode 100644 index 0000000..1e8942b --- /dev/null +++ b/rc_auth/src/robocraft/intercom/status.rs @@ -0,0 +1,19 @@ +use actix_web::{HttpRequest, HttpResponse, web::{Data, Path, Json}, Error, get, post}; + +#[get("/intercom/.status")] +pub async fn status_get(reg: Data) -> Json { + Json(oj_serdes::Status { + servers: reg.statuses().await + }) +} + +#[post("/intercom/.status/{name}/{service}")] +pub async fn status_set(req: HttpRequest, body: Json, auth: Data, reg: Data, uri: Path<(String, String)>) -> Result { + let name = &uri.0; + let service = &uri.1; + log::debug!("Got intercom status message from {}/{}", name, service); + auth.validate(&req, &format!(".status/{}/{}", name, service))?; + log::debug!("Authenticated intercom status message from {}/{}", name, service); + reg.save_status(service.clone(), body.clone()).await; + Ok(HttpResponse::NoContent().finish()) +} diff --git a/rc_auth/src/robocraft/intercom/user_registry.rs b/rc_auth/src/robocraft/intercom/user_registry.rs index eac6ab4..831b25f 100644 --- a/rc_auth/src/robocraft/intercom/user_registry.rs +++ b/rc_auth/src/robocraft/intercom/user_registry.rs @@ -1,11 +1,13 @@ pub struct Users { service_listeners: tokio::sync::RwLock>>, + service_status: tokio::sync::RwLock>, } impl Users { pub fn new() -> Self { Self { service_listeners: tokio::sync::RwLock::new(std::collections::HashMap::with_capacity(16)), + service_status: tokio::sync::RwLock::new(std::collections::HashMap::with_capacity(16)), } } @@ -47,4 +49,12 @@ impl Users { } } } + + pub async fn statuses(&self) -> std::collections::HashMap { + self.service_status.read().await.clone() + } + + pub async fn save_status(&self, service: String, status: oj_serdes::ServerStatus) -> Option { + self.service_status.write().await.insert(service, status) + } } diff --git a/rc_chat_room/Cargo.toml b/rc_chat_room/Cargo.toml index b81a666..285d131 100644 --- a/rc_chat_room/Cargo.toml +++ b/rc_chat_room/Cargo.toml @@ -22,3 +22,4 @@ regex = "1" async-trait.workspace = true git-version.workspace = true chrono.workspace = true +oj_serdes.workspace = true diff --git a/rc_chat_room/src/main.rs b/rc_chat_room/src/main.rs index 377734d..321d31d 100644 --- a/rc_chat_room/src/main.rs +++ b/rc_chat_room/src/main.rs @@ -15,10 +15,11 @@ use tokio::net; use polariton::packet::{Data, Message, Packet, StandardMessage}; use polariton::operation::{OperationResponse, Typed}; -pub type UserTy = oj_rc_core::UserState; +pub type UserTy = std::sync::Arc; pub static START_TIMESTAMP_S: std::sync::atomic::AtomicI64 = std::sync::atomic::AtomicI64::new(0); pub static READY_DURATION_NS: std::sync::atomic::AtomicI64 = std::sync::atomic::AtomicI64::new(0); +pub static ONLINE_USERS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); #[tokio::main] async fn main() -> std::io::Result<()> { @@ -66,11 +67,16 @@ async fn process_socket(mut socket: net::TcpStream, address: std::net::SocketAdd return; } }; + ONLINE_USERS.fetch_add(1, std::sync::atomic::Ordering::SeqCst); let (chann_tx, chann_rx) = tokio::sync::mpsc::unbounded_channel(); - let user_state = oj_rc_core::UserState::<()>::new(users, chann_tx.clone()); + let user_state = std::sync::Arc::new(oj_rc_core::UserState::<()>::new(users, chann_tx.clone())); let (socket_r, socket_w) = socket.into_split(); - server.handle_async_with_channel(socket_r, socket_w, user_state, polariton::packet::SerdesContext::from_boxed(Default::default(), enc), chann_tx, chann_rx).await; + server.handle_async_with_channel_join(socket_r, socket_w, user_state.clone(), polariton::packet::SerdesContext::from_boxed(Default::default(), enc), chann_tx, chann_rx).await; log::debug!("Goodbye connection from address {}", address); + ONLINE_USERS.fetch_sub(1, std::sync::atomic::Ordering::SeqCst); + if let Ok(user_info) = user_state.user() { + update_status(user_info.as_ref().as_ref()).await; + } } const APP_ID: &str = "ChatServer"; @@ -276,3 +282,14 @@ async fn do_connect_handshake( Some(ctx.into_crypto()) } + +pub async fn update_status(user_info: &dyn oj_rc_core::persist::user::IntercomUser) { + user_info.update_status( + env!("CARGO_PKG_NAME"), + oj_serdes::ServerStatus { + uptime_s: (chrono::Utc::now().timestamp() - crate::START_TIMESTAMP_S.load(std::sync::atomic::Ordering::Relaxed)).try_into().unwrap_or_default(), + players: ONLINE_USERS.load(std::sync::atomic::Ordering::SeqCst), + version: env!("CARGO_PKG_VERSION").to_owned(), + }, + ).await; +} diff --git a/rc_chat_room/src/operations/more_auth.rs b/rc_chat_room/src/operations/more_auth.rs index 2077f7f..135e3eb 100644 --- a/rc_chat_room/src/operations/more_auth.rs +++ b/rc_chat_room/src/operations/more_auth.rs @@ -36,6 +36,7 @@ impl MoreLobbyAuth { //let channels = user_impl.subscribed_channels_strings().await?; //let event_tx = user.event_chann(); //self.chat_system.system_mut().connect_user(name, channels, event_tx); + crate::update_status(user.user().unwrap().as_ref().as_ref()).await; let mut resp_params = std::collections::HashMap::new(); resp_params.insert(Self::AUTH_PAYLOAD_KEY, polariton::operation::Typed::Byte(0)); return Ok(resp_params.into()); diff --git a/rc_core/Cargo.toml b/rc_core/Cargo.toml index f9303d7..dddc3f9 100644 --- a/rc_core/Cargo.toml +++ b/rc_core/Cargo.toml @@ -31,6 +31,7 @@ argon2 = { version = "0.5", features = [ "std" ] } sha2 = "0.10" reqwest = { version = "0.12", default-features = false, features = [ "rustls-tls", "charset" ] } reqwest-websocket = { version = "0.5", default-features = false, features = [ "json" ] } +oj_serdes.workspace = true oj_rc_database = { version = "*", path = "../rc_database" } oj_rc_factory = { version = "*", path = "../rc_factory" } diff --git a/rc_core/src/persist/user/intercom.rs b/rc_core/src/persist/user/intercom.rs index 8157be3..643c9df 100644 --- a/rc_core/src/persist/user/intercom.rs +++ b/rc_core/src/persist/user/intercom.rs @@ -83,6 +83,12 @@ impl super::IntercomUser for super::account_json::UserData { log::error!("Failed to send intercom maintenance mode message: {}", e); } } + + async fn update_status(&self, server_name: &str, msg: oj_serdes::ServerStatus) { + if let Err(e) = self.post_to_intercom(&msg, ".status", server_name).await { + log::error!("Failed to send intercom status message: {}", e); + } + } } #[derive(Serialize, Deserialize, Clone, Debug)] diff --git a/rc_core/src/persist/user/traits.rs b/rc_core/src/persist/user/traits.rs index fdb36b5..aa0073d 100644 --- a/rc_core/src/persist/user/traits.rs +++ b/rc_core/src/persist/user/traits.rs @@ -323,7 +323,7 @@ pub enum MultiplayerErrorCode { } #[async_trait::async_trait] -pub trait MultiplayerUser: CommonUser { +pub trait MultiplayerUser: IntercomUser + CommonUser { fn user_id(&self) -> i32; fn user_name(&self) -> &'_ str; fn display_name(&self) -> &'_ str; @@ -334,11 +334,12 @@ pub trait MultiplayerUser: CommonUser { } #[async_trait::async_trait] -pub trait IntercomUser { +pub trait IntercomUser: CommonUser { async fn save_custom_avatar(&self, image: Vec) -> 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 update_status(&self, server_name: &str, msg: oj_serdes::ServerStatus); } pub struct IntercomListener { diff --git a/rc_lobby_room/Cargo.toml b/rc_lobby_room/Cargo.toml index 557ea10..4e1aa12 100644 --- a/rc_lobby_room/Cargo.toml +++ b/rc_lobby_room/Cargo.toml @@ -19,3 +19,4 @@ oj_rc_core = { version = "*", path = "../rc_core" } oj_rc_factory = { version = "*", path = "../rc_factory" } async-trait.workspace = true chrono.workspace = true +oj_serdes.workspace = true diff --git a/rc_lobby_room/src/main.rs b/rc_lobby_room/src/main.rs index b809761..63a2768 100644 --- a/rc_lobby_room/src/main.rs +++ b/rc_lobby_room/src/main.rs @@ -22,6 +22,9 @@ pub struct InitConfig { pub type UserTy = std::sync::Arc>; +pub static START_TIMESTAMP_S: std::sync::atomic::AtomicI64 = std::sync::atomic::AtomicI64::new(0); +pub static ONLINE_USERS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + #[tokio::main] async fn main() -> std::io::Result<()> { env_logger::init(); @@ -47,6 +50,9 @@ async fn main() -> std::io::Result<()> { let listener = net::TcpListener::bind(std::net::SocketAddr::new(ip_addr, args.port)).await?; + let start_time = chrono::Utc::now(); + START_TIMESTAMP_S.store(start_time.timestamp(), std::sync::atomic::Ordering::Relaxed); + if args.once { log::warn!("Handling first connection and then exiting"); let (socket, address) = listener.accept().await?; @@ -71,6 +77,7 @@ async fn process_socket(mut socket: net::TcpStream, address: std::net::SocketAdd return; } }; + ONLINE_USERS.fetch_add(1, std::sync::atomic::Ordering::SeqCst); let (socket_r, socket_w) = socket.into_split(); let (chann_tx, chann_rx) = tokio::sync::mpsc::unbounded_channel(); let user_state = std::sync::Arc::new(oj_rc_core::UserState::<()>::new(users, chann_tx.clone())); @@ -82,6 +89,10 @@ async fn process_socket(mut socket: net::TcpStream, address: std::net::SocketAdd log::debug!("Unauthenticated user disconnected"); } log::debug!("Goodbye connection from address {}", address); + ONLINE_USERS.fetch_sub(1, std::sync::atomic::Ordering::SeqCst); + if let Ok(user_info) = user_state.user() { + update_status(user_info.as_ref().as_ref()).await; + } } const APP_ID: &str = "LobbyServer"; @@ -287,3 +298,14 @@ async fn do_connect_handshake( Some(ctx.into_crypto()) } + +pub async fn update_status(user_info: &dyn oj_rc_core::persist::user::IntercomUser) { + user_info.update_status( + env!("CARGO_PKG_NAME"), + oj_serdes::ServerStatus { + uptime_s: (chrono::Utc::now().timestamp() - crate::START_TIMESTAMP_S.load(std::sync::atomic::Ordering::Relaxed)).try_into().unwrap_or_default(), + players: ONLINE_USERS.load(std::sync::atomic::Ordering::SeqCst), + version: env!("CARGO_PKG_VERSION").to_owned(), + }, + ).await; +} diff --git a/rc_lobby_room/src/operations/more_auth.rs b/rc_lobby_room/src/operations/more_auth.rs index 6b0f69f..b5d1730 100644 --- a/rc_lobby_room/src/operations/more_auth.rs +++ b/rc_lobby_room/src/operations/more_auth.rs @@ -16,6 +16,7 @@ impl Operation for MoreLobbyAuth { if let Some(Typed::Str(auth_payload)) = params_dict.get(&Self::AUTH_PAYLOAD_KEY) { //let mut write_lock = user.write().unwrap(); if user.update_with_auth(&auth_payload.string).await { + crate::update_status(user.user().unwrap().as_ref().as_ref()).await; let mut resp_params = std::collections::HashMap::new(); resp_params.insert(Self::AUTH_PAYLOAD_KEY, polariton::operation::Typed::Byte(0)); return polariton::operation::OperationResponse { diff --git a/rc_multiplayer/Cargo.toml b/rc_multiplayer/Cargo.toml index 9c223f3..3794127 100644 --- a/rc_multiplayer/Cargo.toml +++ b/rc_multiplayer/Cargo.toml @@ -26,3 +26,5 @@ rand.workspace = true literustlib_server = { version = "0.2" } literustlib = { version = "0.2" } rlnl = { version = "0.1", path = "../../rlnl" } + +oj_serdes.workspace = true diff --git a/rc_multiplayer/src/main.rs b/rc_multiplayer/src/main.rs index 4fe1053..a26029d 100644 --- a/rc_multiplayer/src/main.rs +++ b/rc_multiplayer/src/main.rs @@ -16,6 +16,8 @@ pub struct InitConfig { pub matches_chann: tokio::sync::mpsc::Sender, } +pub static START_TIMESTAMP_S: std::sync::atomic::AtomicI64 = std::sync::atomic::AtomicI64::new(0); + #[tokio::main] async fn main() -> std::io::Result<()> { env_logger::init(); @@ -41,5 +43,19 @@ async fn main() -> std::io::Result<()> { 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"); + let start_time = chrono::Utc::now(); + START_TIMESTAMP_S.store(start_time.timestamp(), std::sync::atomic::Ordering::Relaxed); + server.listen().await } + +pub async fn update_status(user_info: &dyn oj_rc_core::persist::user::IntercomUser, player_count: u64) { + user_info.update_status( + env!("CARGO_PKG_NAME"), + oj_serdes::ServerStatus { + uptime_s: (chrono::Utc::now().timestamp() - crate::START_TIMESTAMP_S.load(std::sync::atomic::Ordering::Relaxed)).try_into().unwrap_or_default(), + players: player_count, + version: env!("CARGO_PKG_VERSION").to_owned(), + }, + ).await; +} diff --git a/rc_multiplayer/src/matches/aggregate.rs b/rc_multiplayer/src/matches/aggregate.rs index 60f927d..2fec6e8 100644 --- a/rc_multiplayer/src/matches/aggregate.rs +++ b/rc_multiplayer/src/matches/aggregate.rs @@ -220,17 +220,18 @@ impl GameMatches { if let Some(tx) = self.matches.get(&game_guid) { if tx.is_closed() { self.do_game_cleanup(&game_guid); - self.create_new_game(user, game_guid, connection, response, sender).await; + self.create_new_game(user.clone(), game_guid, connection, response, sender).await; } else { self.routing.insert(user.user_id(), game_guid.clone()); - if tx.send(super::GameMessage::NewConnection { user, game_guid, connection, response, sender }).await.is_err() { + if tx.send(super::GameMessage::NewConnection { user: user.clone(), game_guid, connection, response, sender }).await.is_err() { log::error!("Failed to send NewConnection game message to existing match"); } } } else { - self.create_new_game(user, game_guid, connection, response, sender).await; + self.create_new_game(user.clone(), game_guid, connection, response, sender).await; } - } + crate::update_status(user.as_ref().as_ref(), self.routing.len() as u64).await; + }, msg => { let user_id = msg.user_id(); let mut to_clean = None; diff --git a/rc_services_room/Cargo.toml b/rc_services_room/Cargo.toml index 4050050..1c5c8b4 100644 --- a/rc_services_room/Cargo.toml +++ b/rc_services_room/Cargo.toml @@ -25,3 +25,4 @@ oj_rc_factory = { version = "*", path = "../rc_factory" } rand.workspace = true async-trait.workspace = true libfj.workspace = true +oj_serdes.workspace = true diff --git a/rc_services_room/src/main.rs b/rc_services_room/src/main.rs index ee699a0..00edef9 100644 --- a/rc_services_room/src/main.rs +++ b/rc_services_room/src/main.rs @@ -11,7 +11,10 @@ use tokio::net; use polariton::packet::{Data, Message, Packet, StandardMessage}; use polariton::operation::{OperationResponse, Typed}; -pub type UserTy = oj_rc_core::UserState<()>; +pub type UserTy = std::sync::Arc>; + +pub static START_TIMESTAMP_S: std::sync::atomic::AtomicI64 = std::sync::atomic::AtomicI64::new(0); +pub static ONLINE_USERS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); pub struct InitConfig { pub cubes: oj_rc_core::persist::config::ConfigImpl, @@ -43,6 +46,9 @@ async fn main() -> std::io::Result<()> { let listener = net::TcpListener::bind(std::net::SocketAddr::new(ip_addr, args.port)).await?; + let start_time = chrono::Utc::now(); + START_TIMESTAMP_S.store(start_time.timestamp(), std::sync::atomic::Ordering::Relaxed); + if args.once { log::warn!("Handling first connection and then exiting"); let (socket, address) = listener.accept().await?; @@ -67,12 +73,17 @@ async fn process_socket(mut socket: net::TcpStream, address: std::net::SocketAdd return; } }; + ONLINE_USERS.fetch_add(1, std::sync::atomic::Ordering::SeqCst); let (socket_r, socket_w) = socket.into_split(); let (chann_tx, chann_rx) = tokio::sync::mpsc::unbounded_channel(); - let user_state = oj_rc_core::UserState::<()>::new(init_ctx.users.clone(), chann_tx.clone()); + let user_state = std::sync::Arc::new(oj_rc_core::UserState::<()>::new(init_ctx.users.clone(), chann_tx.clone())); let ctx = polariton::packet::SerdesContext::from_boxed(Default::default(), enc); - server.handle_async_with_channel(socket_r, socket_w, user_state, ctx, chann_tx, chann_rx).await; + server.handle_async_with_channel_join(socket_r, socket_w, user_state.clone(), ctx, chann_tx, chann_rx).await; log::debug!("Goodbye connection from address {}", address); + ONLINE_USERS.fetch_sub(1, std::sync::atomic::Ordering::SeqCst); + if let Ok(user_info) = user_state.user() { + update_status(user_info.as_ref().as_ref()).await; + } } const APP_ID: &str = "WebServicesServer"; @@ -278,3 +289,14 @@ async fn do_connect_handshake( Some(ctx.into_crypto()) } + +pub async fn update_status(user_info: &dyn oj_rc_core::persist::user::IntercomUser) { + user_info.update_status( + env!("CARGO_PKG_NAME"), + oj_serdes::ServerStatus { + uptime_s: (chrono::Utc::now().timestamp() - crate::START_TIMESTAMP_S.load(std::sync::atomic::Ordering::Relaxed)).try_into().unwrap_or_default(), + players: ONLINE_USERS.load(std::sync::atomic::Ordering::SeqCst), + version: env!("CARGO_PKG_VERSION").to_owned(), + }, + ).await; +} diff --git a/rc_services_room/src/operations/more_auth.rs b/rc_services_room/src/operations/more_auth.rs index 27f3389..8b2b014 100644 --- a/rc_services_room/src/operations/more_auth.rs +++ b/rc_services_room/src/operations/more_auth.rs @@ -27,6 +27,7 @@ impl Operation for MoreLobbyAuth { } else { match user_info.webservice_listener().await { Ok(listener) => { + crate::update_status(user_info.as_ref().as_ref()).await; let mut resp_params = std::collections::HashMap::with_capacity(1); resp_params.insert(Self::AUTH_PAYLOAD_KEY, polariton::operation::Typed::Byte(0)); crate::events::IntercomHandler::new(listener, &user_info, user.event_sender()).run(); diff --git a/rc_singleplayer_room/Cargo.toml b/rc_singleplayer_room/Cargo.toml index 459aa07..fbf863a 100644 --- a/rc_singleplayer_room/Cargo.toml +++ b/rc_singleplayer_room/Cargo.toml @@ -17,3 +17,5 @@ oj_polariton_auth = { version = "*", path = "../polariton_auth" } polariton_server.workspace = true oj_rc_core = { version = "*", path = "../rc_core" } async-trait.workspace = true +chrono.workspace = true +oj_serdes.workspace = true diff --git a/rc_singleplayer_room/src/main.rs b/rc_singleplayer_room/src/main.rs index d52129c..92646c7 100644 --- a/rc_singleplayer_room/src/main.rs +++ b/rc_singleplayer_room/src/main.rs @@ -17,7 +17,10 @@ pub struct InitConfig { pub parsers: oj_rc_core::cubes::CubeParsers, } -pub type UserTy = oj_rc_core::UserState<()>; +pub type UserTy = std::sync::Arc>; + +pub static START_TIMESTAMP_S: std::sync::atomic::AtomicI64 = std::sync::atomic::AtomicI64::new(0); +pub static ONLINE_USERS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); #[tokio::main] async fn main() -> std::io::Result<()> { @@ -43,6 +46,9 @@ async fn main() -> std::io::Result<()> { let listener = net::TcpListener::bind(std::net::SocketAddr::new(ip_addr, args.port)).await?; + let start_time = chrono::Utc::now(); + START_TIMESTAMP_S.store(start_time.timestamp(), std::sync::atomic::Ordering::Relaxed); + if args.once { log::warn!("Handling first connection and then exiting"); let (socket, address) = listener.accept().await?; @@ -67,12 +73,17 @@ async fn process_socket(mut socket: net::TcpStream, address: std::net::SocketAdd return; } }; + ONLINE_USERS.fetch_add(1, std::sync::atomic::Ordering::SeqCst); let (socket_r, socket_w) = socket.into_split(); let (chann_tx, chann_rx) = tokio::sync::mpsc::unbounded_channel(); - let user_state = oj_rc_core::UserState::<()>::new(users, chann_tx.clone()); + let user_state = std::sync::Arc::new(oj_rc_core::UserState::<()>::new(users, chann_tx.clone())); let ctx = polariton::packet::SerdesContext::from_boxed(Default::default(), enc); - server.handle_async_with_channel(socket_r, socket_w, user_state, ctx, chann_tx, chann_rx).await; + server.handle_async_with_channel_join(socket_r, socket_w, user_state.clone(), ctx, chann_tx, chann_rx).await; log::debug!("Goodbye connection from address {}", address); + ONLINE_USERS.fetch_sub(1, std::sync::atomic::Ordering::SeqCst); + if let Ok(user_info) = user_state.user() { + update_status(user_info.as_ref().as_ref()).await; + } } const APP_ID: &str = "SinglePlayerServer"; @@ -278,3 +289,14 @@ async fn do_connect_handshake( Some(ctx.into_crypto()) } + +pub async fn update_status(user_info: &dyn oj_rc_core::persist::user::IntercomUser) { + user_info.update_status( + env!("CARGO_PKG_NAME"), + oj_serdes::ServerStatus { + uptime_s: (chrono::Utc::now().timestamp() - crate::START_TIMESTAMP_S.load(std::sync::atomic::Ordering::Relaxed)).try_into().unwrap_or_default(), + players: ONLINE_USERS.load(std::sync::atomic::Ordering::SeqCst), + version: env!("CARGO_PKG_VERSION").to_owned(), + }, + ).await; +} diff --git a/rc_singleplayer_room/src/operations/more_auth.rs b/rc_singleplayer_room/src/operations/more_auth.rs index 6b0f69f..b5d1730 100644 --- a/rc_singleplayer_room/src/operations/more_auth.rs +++ b/rc_singleplayer_room/src/operations/more_auth.rs @@ -16,6 +16,7 @@ impl Operation for MoreLobbyAuth { if let Some(Typed::Str(auth_payload)) = params_dict.get(&Self::AUTH_PAYLOAD_KEY) { //let mut write_lock = user.write().unwrap(); if user.update_with_auth(&auth_payload.string).await { + crate::update_status(user.user().unwrap().as_ref().as_ref()).await; let mut resp_params = std::collections::HashMap::new(); resp_params.insert(Self::AUTH_PAYLOAD_KEY, polariton::operation::Typed::Byte(0)); return polariton::operation::OperationResponse { diff --git a/rc_social_room/Cargo.toml b/rc_social_room/Cargo.toml index ff9b392..38115fe 100644 --- a/rc_social_room/Cargo.toml +++ b/rc_social_room/Cargo.toml @@ -17,3 +17,5 @@ oj_polariton_auth = { version = "*", path = "../polariton_auth" } polariton_server.workspace = true oj_rc_core = { version = "*", path = "../rc_core" } async-trait.workspace = true +chrono.workspace = true +oj_serdes.workspace = true diff --git a/rc_social_room/src/main.rs b/rc_social_room/src/main.rs index 455cd75..1c109f3 100644 --- a/rc_social_room/src/main.rs +++ b/rc_social_room/src/main.rs @@ -10,7 +10,10 @@ use tokio::net; use polariton::packet::{Data, Message, Packet, StandardMessage}; use polariton::operation::{OperationResponse, Typed}; -pub type UserTy = oj_rc_core::UserState; +pub type UserTy = std::sync::Arc>; + +pub static START_TIMESTAMP_S: std::sync::atomic::AtomicI64 = std::sync::atomic::AtomicI64::new(0); +pub static ONLINE_USERS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); #[tokio::main] async fn main() -> std::io::Result<()> { @@ -27,6 +30,9 @@ async fn main() -> std::io::Result<()> { let listener = net::TcpListener::bind(std::net::SocketAddr::new(ip_addr, args.port)).await?; + let start_time = chrono::Utc::now(); + START_TIMESTAMP_S.store(start_time.timestamp(), std::sync::atomic::Ordering::Relaxed); + if args.once { log::warn!("Handling first connection and then exiting"); let (socket, address) = listener.accept().await?; @@ -51,13 +57,18 @@ async fn process_socket(mut socket: net::TcpStream, address: std::net::SocketAdd return; } }; + ONLINE_USERS.fetch_add(1, std::sync::atomic::Ordering::SeqCst); let (socket_r, socket_w) = socket.into_split(); let (chann_tx, chann_rx) = tokio::sync::mpsc::unbounded_channel(); - let user_state = oj_rc_core::UserState::::new(users, chann_tx.clone()); + let user_state = std::sync::Arc::new(oj_rc_core::UserState::::new(users, chann_tx.clone())); let op_ctx = polariton::serdes::SerdesContext::::default_const(); let ctx = polariton::packet::SerdesContext::from_boxed(op_ctx, enc); - server.handle_async_with_channel(socket_r, socket_w, user_state, ctx, chann_tx, chann_rx).await; + server.handle_async_with_channel_join(socket_r, socket_w, user_state.clone(), ctx, chann_tx, chann_rx).await; log::debug!("Goodbye connection from address {}", address); + ONLINE_USERS.fetch_sub(1, std::sync::atomic::Ordering::SeqCst); + if let Ok(user_info) = user_state.user() { + update_status(user_info.as_ref().as_ref()).await; + } } const APP_ID: &str = "SocialServer"; @@ -263,3 +274,14 @@ async fn do_connect_handshake( Some(ctx.into_crypto()) } + +pub async fn update_status(user_info: &dyn oj_rc_core::persist::user::IntercomUser) { + user_info.update_status( + env!("CARGO_PKG_NAME"), + oj_serdes::ServerStatus { + uptime_s: (chrono::Utc::now().timestamp() - crate::START_TIMESTAMP_S.load(std::sync::atomic::Ordering::Relaxed)).try_into().unwrap_or_default(), + players: ONLINE_USERS.load(std::sync::atomic::Ordering::SeqCst), + version: env!("CARGO_PKG_VERSION").to_owned(), + }, + ).await; +} diff --git a/rc_social_room/src/operations/more_auth.rs b/rc_social_room/src/operations/more_auth.rs index 47d515e..54a08bd 100644 --- a/rc_social_room/src/operations/more_auth.rs +++ b/rc_social_room/src/operations/more_auth.rs @@ -15,6 +15,7 @@ impl Operation for MoreLobbyAuth { let params_dict = params.to_dict(); if let Some(Typed::Str(auth_payload)) = params_dict.get(&Self::AUTH_PAYLOAD_KEY) { if user.update_with_auth(&auth_payload.string).await { + crate::update_status(user.user().unwrap().as_ref().as_ref()).await; let mut resp_params = std::collections::HashMap::with_capacity(1); resp_params.insert(Self::AUTH_PAYLOAD_KEY, polariton::operation::Typed::Byte(0)); return polariton::operation::OperationResponse {