1
0
mirror of https://git.ngram.ca/OpenJam/rc-servers synced 2026-08-23 23:08:52 +00:00

Implement lobby parts of custom games #38

This commit is contained in:
NG (Graham)
2026-03-17 21:52:15 -04:00
parent af04c4093c
commit 6a953bc613
32 changed files with 1045 additions and 61 deletions

View File

@@ -0,0 +1,63 @@
use actix_web::{rt, web::{Payload, Data, Json, Path}, Error, HttpRequest, HttpResponse, get, post};
//use actix_ws::AggregatedMessage;
//use futures::StreamExt as _;
#[get("/intercom/userless/.oj_lobby")]
pub async fn lobby_state_ws(req: HttpRequest, stream: Payload, auth: Data<super::IntercomAuth>, reg: Data<super::Users>) -> Result<HttpResponse, Error> {
auth.validate(&req, "state/.oj_lobby")?;
let (res, mut session, _stream) = actix_ws::handle(&req, 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_lobby_state_service(tx).await;
log::debug!("Registered lobby state intercom websocket");
// start task but don't wait for it
rt::spawn(async move {
let mut is_ok = false;
while let Some(op) = rx.recv().await {
match op {
super::IntercomOp::Message(msg) => {
if let Err(e) = session.text(serde_json::to_string(&msg).unwrap()).await {
log::warn!("Failed to send lobby state intercom message: {}", e);
break;
}
},
super::IntercomOp::Info(info) => {
match info {
super::IntercomInfo::Close => {
is_ok = true;
break;
},
}
}
}
}
if !is_ok {
reg.remove_lobby_state_service().await;
}
rx.close();
session.close(Some(actix_ws::CloseReason {
code: actix_ws::CloseCode::Normal,
description: Some("End of channel".to_owned()),
})).await.expect("Failed to close a services intercom websocket session");
log::debug!("Lobby state intercom websocket closed");
});
// respond immediately with response connected to WS session
Ok(res)
}
#[post("/intercom/.oj_lobby/{name}/state")]
pub async fn lobby_state_msg(req: HttpRequest, body: Json<oj_rc_core::persist::user::intercom::IntercomLobbyStateMessage>, auth: Data<super::IntercomAuth>, reg: Data<super::Users>, name: Path<String>) -> Result<HttpResponse, super::IntercomOpError> {
log::debug!("Got lobby state intercom message");
auth.validate(&req, &format!(".oj_lobby/{}/state", name))?;
log::debug!("Authenticated intercom lobby state message from {}", name);
reg.broadcast_lobby_message(body.0).await;
Ok(HttpResponse::NoContent().finish())
}

View File

@@ -10,11 +10,17 @@ pub use user_registry::Users;
mod status;
pub use status::{status_set, status_get};
enum IntercomOp {
Message(oj_rc_core::persist::user::intercom::IntercomWebServiceUserMessage),
mod lobby;
pub use lobby::{lobby_state_ws, lobby_state_msg};
enum IntercomOp<M: serde::Serialize + serde::Deserialize<'static> + 'static> {
Message(M),
Info(IntercomInfo),
}
enum IntercomInfo {
Close,
}
type WebServicesIntercomOp = IntercomOp<oj_rc_core::persist::user::intercom::IntercomWebServiceUserMessage>;
type LobbyStateIntercomOp = IntercomOp<oj_rc_core::persist::user::intercom::IntercomLobbyStateMessage>;

View File

@@ -13,7 +13,7 @@ pub async fn services_ws(req: HttpRequest, stream: Payload, auth: Data<super::In
*/
let (tx, mut rx) = tokio::sync::mpsc::channel(16);
reg.register_service(name.clone(), tx).await;
reg.register_web_service(name.clone(), tx).await;
log::debug!("Registered web services intercom websocket for user {}", name);
// start task but don't wait for it
@@ -39,7 +39,7 @@ pub async fn services_ws(req: HttpRequest, stream: Payload, auth: Data<super::In
}
if !is_ok {
reg.remove_service(name.clone()).await;
reg.remove_web_service(name.clone()).await;
}
rx.close();
session.close(Some(actix_ws::CloseReason {
@@ -58,6 +58,6 @@ pub async fn service_msg(req: HttpRequest, body: Json<oj_rc_core::persist::user:
log::debug!("Got intercom message from {} to {:?}", name, body.public_ids.as_slice());
auth.validate(&req, &format!(".oj_services/{}/messages", name))?;
log::debug!("Authenticated intercom message from {} to {:?}", name, body.public_ids.as_slice());
reg.broadcast_service_message(body.0).await;
reg.broadcast_web_service_message(body.0).await;
Ok(HttpResponse::NoContent().finish())
}

View File

@@ -1,18 +1,26 @@
struct ServiceListeners {
web_services: tokio::sync::RwLock<std::collections::HashMap<String, tokio::sync::mpsc::Sender<super::WebServicesIntercomOp>>>,
lobby_state: tokio::sync::RwLock<Option<tokio::sync::mpsc::Sender<super::LobbyStateIntercomOp>>>,
}
pub struct Users {
service_listeners: tokio::sync::RwLock<std::collections::HashMap<String, tokio::sync::mpsc::Sender<super::IntercomOp>>>,
service_listeners: ServiceListeners,
service_status: tokio::sync::RwLock<std::collections::HashMap<String, oj_serdes::ServerStatus>>,
}
impl Users {
pub fn new() -> Self {
Self {
service_listeners: tokio::sync::RwLock::new(std::collections::HashMap::with_capacity(16)),
service_listeners: ServiceListeners {
web_services: tokio::sync::RwLock::new(std::collections::HashMap::with_capacity(16)),
lobby_state: tokio::sync::RwLock::new(None),
},
service_status: tokio::sync::RwLock::new(std::collections::HashMap::with_capacity(16)),
}
}
pub(super) async fn register_service(&self, public_id: String, sender: tokio::sync::mpsc::Sender<super::IntercomOp>) {
let mut write_lock = self.service_listeners.write().await;
pub(super) async fn register_web_service(&self, public_id: String, sender: tokio::sync::mpsc::Sender<super::WebServicesIntercomOp>) {
let mut write_lock = self.service_listeners.web_services.write().await;
if let Some(old_sender) = write_lock.insert(public_id.clone(), sender) {
if !old_sender.is_closed() {
log::warn!("Replaced web services intercom channel for user {} (why duplicate!?)", public_id);
@@ -21,26 +29,43 @@ impl Users {
}
}
pub async fn remove_service(&self, public_id: String) {
let mut write_lock = self.service_listeners.write().await;
pub(super) async fn register_lobby_state_service(&self, sender: tokio::sync::mpsc::Sender<super::LobbyStateIntercomOp>) {
let mut write_lock = self.service_listeners.lobby_state.write().await;
if let Some(old_sender) = write_lock.replace(sender) {
if !old_sender.is_closed() {
log::warn!("Replaced lobby state intercom channel (why duplicate!?)");
old_sender.send(super::IntercomOp::Info(super::IntercomInfo::Close)).await.unwrap_or_default()
}
}
}
pub async fn remove_web_service(&self, public_id: String) {
let mut write_lock = self.service_listeners.web_services.write().await;
if write_lock.remove(&public_id).is_none() {
log::warn!("Tried to remove web services intercom channel for user {} without listener", public_id);
}
}
pub async fn broadcast_service_message(&self, msg: oj_rc_core::persist::user::intercom::IntercomWebServiceMessage) {
let read_lock = self.service_listeners.read().await;
pub async fn remove_lobby_state_service(&self) {
let mut write_lock = self.service_listeners.lobby_state.write().await;
if write_lock.take().is_none() {
log::warn!("Tried to remove lobby state intercom channel without listener");
}
}
pub async fn broadcast_web_service_message(&self, msg: oj_rc_core::persist::user::intercom::IntercomWebServiceMessage) {
let read_lock = self.service_listeners.web_services.read().await;
if msg.everyone {
if !msg.public_ids.is_empty() { return; } // invalid
for (public_id, tx) in read_lock.iter() {
if let Err(e) = tx.send(super::IntercomOp::Message(msg.data.clone())).await {
if let Err(e) = tx.send(super::WebServicesIntercomOp::Message(msg.data.clone())).await {
log::error!("Failed to send web service intercom message to {}: {}", public_id, e);
}
}
} else {
for public_id in msg.public_ids {
if let Some(tx) = read_lock.get(&public_id) {
if let Err(e) = tx.send(super::IntercomOp::Message(msg.data.clone())).await {
if let Err(e) = tx.send(super::WebServicesIntercomOp::Message(msg.data.clone())).await {
log::error!("Failed to send web service intercom message to {}: {}", public_id, e);
}
} else {
@@ -50,6 +75,15 @@ impl Users {
}
}
pub async fn broadcast_lobby_message(&self, msg: oj_rc_core::persist::user::intercom::IntercomLobbyStateMessage) {
let read_lock = self.service_listeners.lobby_state.read().await;
if let Some(listener) = &*read_lock {
if let Err(e) = listener.send(super::LobbyStateIntercomOp::Message(msg)).await {
log::error!("Failed to send lobby state intercom message: {}", e);
}
}
}
pub async fn statuses(&self) -> std::collections::HashMap<String, oj_serdes::ServerStatus> {
self.service_status.read().await.clone()
}