mirror of
https://git.ngram.ca/OpenJam/rc-servers
synced 2026-08-23 23:08:52 +00:00
### Description This PR implements a queue timeout in the lobby server. * After the time configured in `timeout` has elapsed (default: 3 minutes) since the first player joins the queue, matchmaking will start even if the queue has not reached the `players_per_game` target. Any missing slots will be filled with fake players. * If the first queued player leaves and other players remain, the timer is recalculated so that matchmaking starts after the `timeout` duration (default: 3 minutes) from the next player’s join time. (After that, the timer continues to be adjusted in the same way whenever the head of the queue changes.) ### Game Robocraft ### Please confirm - [x] I am the legal owner or represent the owner of all work submitted - [x] I consent to my submission being added to this FOSS project - [x] This PR used LLMs to generate some or all of the code changes Reviewed-on: https://git.ngram.ca/OpenJam/rc-servers/pulls/71 Reviewed-by: NGnius <ngniusness@gmail.com> Co-authored-by: MaxSignal <kastera58@gmail.com> Co-committed-by: MaxSignal <kastera58@gmail.com>
This commit is contained in:
@@ -619,6 +619,7 @@ fn default_multiplayer() -> super::MultiplayerConfig {
|
|||||||
super::MultiplayerConfig {
|
super::MultiplayerConfig {
|
||||||
players_per_game: 2,
|
players_per_game: 2,
|
||||||
enabled: true,
|
enabled: true,
|
||||||
|
autostart_after_s: 180,
|
||||||
network: super::multiplayer::default_net_conf(),
|
network: super::multiplayer::default_net_conf(),
|
||||||
fakes: super::multiplayer::default_fake_users(),
|
fakes: super::multiplayer::default_fake_users(),
|
||||||
battle_arena: super::multiplayer::default_ba_conf(),
|
battle_arena: super::multiplayer::default_ba_conf(),
|
||||||
|
|||||||
@@ -403,6 +403,10 @@ impl <C: Clone + Send> super::ConfigProvider<C> for CubeConfig {
|
|||||||
fn is_multiplayer_enabled(&self) -> bool {
|
fn is_multiplayer_enabled(&self) -> bool {
|
||||||
self.battle.multiplayer.enabled
|
self.battle.multiplayer.enabled
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn multiplayer_autostart_after(&self) -> std::time::Duration {
|
||||||
|
std::time::Duration::from_secs(self.battle.multiplayer.autostart_after_s)
|
||||||
|
}
|
||||||
|
|
||||||
fn network_config(&self) -> crate::persist::NetworkConf {
|
fn network_config(&self) -> crate::persist::NetworkConf {
|
||||||
self.battle.multiplayer.network.clone()
|
self.battle.multiplayer.network.clone()
|
||||||
|
|||||||
@@ -27,6 +27,7 @@ pub trait ConfigProvider<C: Clone> {
|
|||||||
fn singleplayer_details(&self) -> SingleplayerConfig;
|
fn singleplayer_details(&self) -> SingleplayerConfig;
|
||||||
fn players_per_game(&self) -> usize;
|
fn players_per_game(&self) -> usize;
|
||||||
fn is_multiplayer_enabled(&self) -> bool;
|
fn is_multiplayer_enabled(&self) -> bool;
|
||||||
|
fn multiplayer_autostart_after(&self) -> std::time::Duration;
|
||||||
// FIXME don't use serializable types in traits
|
// FIXME don't use serializable types in traits
|
||||||
fn network_config(&self) -> crate::persist::NetworkConf;
|
fn network_config(&self) -> crate::persist::NetworkConf;
|
||||||
fn maps(&self) -> std::collections::HashMap<GameMap, MapConfig>;
|
fn maps(&self) -> std::collections::HashMap<GameMap, MapConfig>;
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ use serde::{Serialize, Deserialize};
|
|||||||
pub struct MultiplayerConfig {
|
pub struct MultiplayerConfig {
|
||||||
pub players_per_game: usize,
|
pub players_per_game: usize,
|
||||||
pub enabled: bool,
|
pub enabled: bool,
|
||||||
|
pub autostart_after_s: u64,
|
||||||
#[serde(default = "default_net_conf")]
|
#[serde(default = "default_net_conf")]
|
||||||
pub network: NetworkConf,
|
pub network: NetworkConf,
|
||||||
#[serde(default = "default_fake_users")]
|
#[serde(default = "default_fake_users")]
|
||||||
|
|||||||
@@ -40,10 +40,12 @@ struct QueueUser {
|
|||||||
emitter: polariton_server::events::EventEmitter,
|
emitter: polariton_server::events::EventEmitter,
|
||||||
player: oj_rc_core::data::player_data::PlayerData,
|
player: oj_rc_core::data::player_data::PlayerData,
|
||||||
user_id: i32,
|
user_id: i32,
|
||||||
|
enqueued_at: chrono::DateTime<chrono::Utc>,
|
||||||
|
user: std::sync::Arc<Box<dyn oj_rc_core::persist::user::User<()> + Send + Sync>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
pub struct QueueHandler {
|
pub struct QueueHandler {
|
||||||
users_in_queue: tokio::sync::Mutex<HashMap<QueueKey, Vec<QueueUser>>>,
|
users_in_queue: std::sync::Arc<tokio::sync::Mutex<HashMap<QueueKey, Vec<QueueUser>>>>,
|
||||||
users_per_game: usize,
|
users_per_game: usize,
|
||||||
is_enabled: bool,
|
is_enabled: bool,
|
||||||
hostname: String,
|
hostname: String,
|
||||||
@@ -53,13 +55,15 @@ pub struct QueueHandler {
|
|||||||
cpu_counter: std::sync::Arc<oj_rc_core::cubes::CpuListParser>,
|
cpu_counter: std::sync::Arc<oj_rc_core::cubes::CpuListParser>,
|
||||||
weapon_guesser: std::sync::Arc<oj_rc_core::cubes::WeaponListParser>,
|
weapon_guesser: std::sync::Arc<oj_rc_core::cubes::WeaponListParser>,
|
||||||
change_strategy: GamemodeChangeStrategy,
|
change_strategy: GamemodeChangeStrategy,
|
||||||
|
autostart_after: std::time::Duration,
|
||||||
|
autostart_task_started: std::sync::atomic::AtomicBool,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl QueueHandler {
|
impl QueueHandler {
|
||||||
pub fn new(conf: &oj_rc_core::ConfigImpl, game_host: &str, factory: std::sync::Arc<oj_rc_core::factory::Factory>, cpu_counter: std::sync::Arc<oj_rc_core::cubes::CpuListParser>, weapon_guesser: std::sync::Arc<oj_rc_core::cubes::WeaponListParser>,) -> Self {
|
pub fn new(conf: &oj_rc_core::ConfigImpl, game_host: &str, factory: std::sync::Arc<oj_rc_core::factory::Factory>, cpu_counter: std::sync::Arc<oj_rc_core::cubes::CpuListParser>, weapon_guesser: std::sync::Arc<oj_rc_core::cubes::WeaponListParser>,) -> Self {
|
||||||
let (domain, port_str) = game_host.split_once(':').expect("Invalid redirect address (must be domain:port)");
|
let (domain, port_str) = game_host.split_once(':').expect("Invalid redirect address (must be domain:port)");
|
||||||
Self {
|
Self {
|
||||||
users_in_queue: tokio::sync::Mutex::new(HashMap::new()),
|
users_in_queue: std::sync::Arc::new(tokio::sync::Mutex::new(HashMap::new())),
|
||||||
users_per_game: oj_rc_core::ConfigProvider::<()>::players_per_game(conf),
|
users_per_game: oj_rc_core::ConfigProvider::<()>::players_per_game(conf),
|
||||||
is_enabled: oj_rc_core::ConfigProvider::<()>::is_multiplayer_enabled(conf),
|
is_enabled: oj_rc_core::ConfigProvider::<()>::is_multiplayer_enabled(conf),
|
||||||
hostname: domain.to_owned(),
|
hostname: domain.to_owned(),
|
||||||
@@ -69,10 +73,122 @@ impl QueueHandler {
|
|||||||
cpu_counter,
|
cpu_counter,
|
||||||
weapon_guesser,
|
weapon_guesser,
|
||||||
change_strategy: GamemodeChangeStrategy::from_core(<oj_rc_core::ConfigImpl as oj_rc_core::ConfigProvider<()>>::server_config(conf).queue_mode),
|
change_strategy: GamemodeChangeStrategy::from_core(<oj_rc_core::ConfigImpl as oj_rc_core::ConfigProvider<()>>::server_config(conf).queue_mode),
|
||||||
|
autostart_after: oj_rc_core::ConfigProvider::<()>::multiplayer_autostart_after(conf),
|
||||||
|
autostart_task_started: std::sync::atomic::AtomicBool::new(false),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn enter_match(&self, key: QueueKey, mut players: Vec<QueueUser>, user: &(dyn oj_rc_core::persist::user::LobbyUser + Send + Sync)) {
|
fn ensure_autostart_task_running(&self) {
|
||||||
|
if self.autostart_task_started.swap(true, std::sync::atomic::Ordering::AcqRel) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
let users_in_queue = self.users_in_queue.clone();
|
||||||
|
let users_per_game = self.users_per_game;
|
||||||
|
let hostname = self.hostname.clone();
|
||||||
|
let hostport = self.hostport;
|
||||||
|
let network_conf = self.network_conf.clone();
|
||||||
|
let factory = self.factory.clone();
|
||||||
|
let cpu_counter = self.cpu_counter.clone();
|
||||||
|
let weapon_guesser = self.weapon_guesser.clone();
|
||||||
|
let autostart_after = self.autostart_after;
|
||||||
|
|
||||||
|
tokio::spawn(async move {
|
||||||
|
loop {
|
||||||
|
let now = chrono::Utc::now();
|
||||||
|
let mut to_start: Vec<(QueueKey, Vec<QueueUser>)> = Vec::new();
|
||||||
|
let mut next_deadline: Option<chrono::DateTime<chrono::Utc>> = None;
|
||||||
|
|
||||||
|
{
|
||||||
|
let mut lock = users_in_queue.lock().await;
|
||||||
|
|
||||||
|
let mut empty_keys: Vec<QueueKey> = Vec::new();
|
||||||
|
let mut expired_keys: Vec<QueueKey> = Vec::new();
|
||||||
|
|
||||||
|
for (key, users) in lock.iter() {
|
||||||
|
if users.is_empty() {
|
||||||
|
empty_keys.push(key.clone());
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
// first user is always oldest within a queue
|
||||||
|
let deadline = users[0].enqueued_at + autostart_after;
|
||||||
|
if now >= deadline {
|
||||||
|
expired_keys.push(key.clone());
|
||||||
|
} else {
|
||||||
|
next_deadline = Some(match next_deadline {
|
||||||
|
Some(d) => if deadline < d { deadline } else { d },
|
||||||
|
None => deadline,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
for k in empty_keys {
|
||||||
|
lock.remove(&k);
|
||||||
|
}
|
||||||
|
|
||||||
|
for k in expired_keys {
|
||||||
|
if let Some(players) = lock.remove(&k) {
|
||||||
|
if !players.is_empty() {
|
||||||
|
to_start.push((k, players));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// start expired queues outside the lock
|
||||||
|
for (key, players) in to_start {
|
||||||
|
// choose the oldest queued user's LobbyUser handle
|
||||||
|
let starter = match players.first() {
|
||||||
|
Some(p) => p.user.clone(),
|
||||||
|
None => continue,
|
||||||
|
};
|
||||||
|
|
||||||
|
QueueHandler::enter_match_static(
|
||||||
|
hostname.clone(),
|
||||||
|
hostport,
|
||||||
|
network_conf.clone(),
|
||||||
|
factory.clone(),
|
||||||
|
cpu_counter.clone(),
|
||||||
|
weapon_guesser.clone(),
|
||||||
|
users_per_game,
|
||||||
|
key,
|
||||||
|
players,
|
||||||
|
starter.as_ref().as_ref(),
|
||||||
|
).await;
|
||||||
|
}
|
||||||
|
|
||||||
|
let sleep_dur = if let Some(deadline) = next_deadline {
|
||||||
|
let now2 = chrono::Utc::now();
|
||||||
|
match (deadline - now2).to_std() {
|
||||||
|
Ok(d) => d,
|
||||||
|
Err(_) => std::time::Duration::from_secs(0),
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
std::time::Duration::from_secs(1)
|
||||||
|
};
|
||||||
|
|
||||||
|
tokio::time::sleep(sleep_dur).await;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn enter_match(&self, key: QueueKey, players: Vec<QueueUser>, user: &(dyn oj_rc_core::persist::user::LobbyUser + Send + Sync)) {
|
||||||
|
Self::enter_match_static(
|
||||||
|
self.hostname.clone(),
|
||||||
|
self.hostport,
|
||||||
|
self.network_conf.clone(),
|
||||||
|
self.factory.clone(),
|
||||||
|
self.cpu_counter.clone(),
|
||||||
|
self.weapon_guesser.clone(),
|
||||||
|
self.users_per_game,
|
||||||
|
key,
|
||||||
|
players,
|
||||||
|
user,
|
||||||
|
).await
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn enter_match_static(hostname: String, hostport: u16, network_conf: crate::data::network::NetworkConfigData, factory: std::sync::Arc<oj_rc_core::factory::Factory>, cpu_counter: std::sync::Arc<oj_rc_core::cubes::CpuListParser>, weapon_guesser: std::sync::Arc<oj_rc_core::cubes::WeaponListParser>, users_per_game: usize, key: QueueKey, mut players: Vec<QueueUser>, user: &(dyn oj_rc_core::persist::user::LobbyUser + Send + Sync)) {
|
||||||
let guid_str = key.guid();
|
let guid_str = key.guid();
|
||||||
let game_desc = oj_rc_core::persist::user::GameDescriptor {
|
let game_desc = oj_rc_core::persist::user::GameDescriptor {
|
||||||
guid: guid_str.clone(),
|
guid: guid_str.clone(),
|
||||||
@@ -100,14 +216,15 @@ impl QueueHandler {
|
|||||||
display_name: x.player.display_name.clone(),
|
display_name: x.player.display_name.clone(),
|
||||||
}).collect();
|
}).collect();
|
||||||
|
|
||||||
match user.start_game(game_desc, player_descs, self.factory.as_ref(), &self.cpu_counter, &self.weapon_guesser, &team_picker).await {
|
match user.start_game(game_desc, player_descs, factory.as_ref(), &cpu_counter, &weapon_guesser, &team_picker).await {
|
||||||
Ok(fakes) => {
|
Ok(fakes) => {
|
||||||
|
let missing = users_per_game.saturating_sub(players.len());
|
||||||
let player_datas = players.iter().map(|x| x.player.clone())
|
let player_datas = players.iter().map(|x| x.player.clone())
|
||||||
.chain(fakes.players.into_iter().map(|(desc, _emu)| desc))
|
.chain(fakes.players.into_iter().take(missing).map(|(desc, _emu)| desc),)
|
||||||
.collect();
|
.collect();
|
||||||
let enter_battle_ev = crate::events::battle_enter::BattleEnter {
|
let enter_battle_ev = crate::events::battle_enter::BattleEnter {
|
||||||
host: self.hostname.clone(),
|
host: hostname.clone(),
|
||||||
port: self.hostport,
|
port: hostport,
|
||||||
map: key.map.clone(),
|
map: key.map.clone(),
|
||||||
mode: key.mode,
|
mode: key.mode,
|
||||||
guid: guid_str.clone(),
|
guid: guid_str.clone(),
|
||||||
@@ -116,7 +233,7 @@ impl QueueHandler {
|
|||||||
visibility: Some(key.visibility),
|
visibility: Some(key.visibility),
|
||||||
auto_heal: key.auto_heal,
|
auto_heal: key.auto_heal,
|
||||||
player_datas,
|
player_datas,
|
||||||
network_config: self.network_conf.clone(),
|
network_config: network_conf.clone(),
|
||||||
};
|
};
|
||||||
let arc_event = std::sync::Arc::new(enter_battle_ev);
|
let arc_event = std::sync::Arc::new(enter_battle_ev);
|
||||||
for player in players.iter() {
|
for player in players.iter() {
|
||||||
@@ -144,7 +261,7 @@ impl QueueHandler {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn join_queue(&self, map: String, mode: oj_rc_core::data::game_mode::GameMode, visibility: oj_rc_core::data::game_mode::MapVisibility, auto_heal: bool, user: &(dyn oj_rc_core::persist::user::LobbyUser + Send + Sync), event_emitter: polariton_server::events::EventEmitter) {
|
pub async fn join_queue(&self, map: String, mode: oj_rc_core::data::game_mode::GameMode, visibility: oj_rc_core::data::game_mode::MapVisibility, auto_heal: bool, user: std::sync::Arc<Box<dyn oj_rc_core::persist::user::User<()> + Send + Sync>>, event_emitter: polariton_server::events::EventEmitter) {
|
||||||
if !self.is_enabled {
|
if !self.is_enabled {
|
||||||
event_emitter.emit(crate::events::enqueue_error::QueueJoinError {
|
event_emitter.emit(crate::events::enqueue_error::QueueJoinError {
|
||||||
code: oj_rc_core::data::error_codes::LobbyReasonCode::NoSuitableLobbyFound as i16,
|
code: oj_rc_core::data::error_codes::LobbyReasonCode::NoSuitableLobbyFound as i16,
|
||||||
@@ -152,15 +269,21 @@ impl QueueHandler {
|
|||||||
});
|
});
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
self.ensure_autostart_task_running();
|
||||||
|
|
||||||
let key = QueueKey {
|
let key = QueueKey {
|
||||||
map, mode, visibility, auto_heal,
|
map, mode, visibility, auto_heal,
|
||||||
};
|
};
|
||||||
match user.player_data(&self.cpu_counter).await {
|
let lobby_user = user.as_ref().as_ref();
|
||||||
|
match lobby_user.player_data(&self.cpu_counter).await {
|
||||||
Ok(player_data) => {
|
Ok(player_data) => {
|
||||||
let new_player = QueueUser {
|
let new_player = QueueUser {
|
||||||
emitter: event_emitter,
|
emitter: event_emitter,
|
||||||
player: player_data,
|
player: player_data,
|
||||||
user_id: user.user_id(),
|
user_id: oj_rc_core::persist::user::LobbyUser::user_id(lobby_user),
|
||||||
|
enqueued_at: chrono::Utc::now(),
|
||||||
|
user: user.clone(),
|
||||||
};
|
};
|
||||||
let mut lock = self.users_in_queue.lock().await;
|
let mut lock = self.users_in_queue.lock().await;
|
||||||
// handle game event change
|
// handle game event change
|
||||||
@@ -177,7 +300,10 @@ impl QueueHandler {
|
|||||||
} else {
|
} else {
|
||||||
new_queue_map.insert(key.clone(), users);
|
new_queue_map.insert(key.clone(), users);
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
for users in new_queue_map.values_mut() {
|
||||||
|
users.sort_by_key(|u| u.enqueued_at);
|
||||||
}
|
}
|
||||||
*lock = new_queue_map;
|
*lock = new_queue_map;
|
||||||
if count != 0 {
|
if count != 0 {
|
||||||
@@ -217,7 +343,7 @@ impl QueueHandler {
|
|||||||
let players = if game_ready { lock.remove(&key) } else { None };
|
let players = if game_ready { lock.remove(&key) } else { None };
|
||||||
drop(lock);
|
drop(lock);
|
||||||
if let Some(players) = players {
|
if let Some(players) = players {
|
||||||
self.enter_match(key, players, user).await;
|
self.enter_match(key, players, lobby_user).await;
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
@@ -229,8 +355,8 @@ impl QueueHandler {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn leave_queue(&self, user: &(dyn oj_rc_core::persist::user::LobbyUser + Send + Sync)) {
|
pub async fn leave_queue(&self, user: std::sync::Arc<Box<dyn oj_rc_core::persist::user::User<()> + Send + Sync>>) {
|
||||||
let user_id = user.user_id();
|
let user_id = oj_rc_core::persist::user::LobbyUser::user_id(user.as_ref().as_ref());
|
||||||
for queue in self.users_in_queue.lock().await.values_mut() {
|
for queue in self.users_in_queue.lock().await.values_mut() {
|
||||||
if let Some((i, _)) = queue.iter().enumerate().find(|(_, user)| user.user_id == user_id) {
|
if let Some((i, _)) = queue.iter().enumerate().find(|(_, user)| user.user_id == user_id) {
|
||||||
queue.remove(i);
|
queue.remove(i);
|
||||||
|
|||||||
@@ -85,7 +85,7 @@ async fn process_socket(mut socket: net::TcpStream, address: std::net::SocketAdd
|
|||||||
let ctx = polariton::packet::SerdesContext::from_boxed(Default::default(), enc);
|
let ctx = polariton::packet::SerdesContext::from_boxed(Default::default(), enc);
|
||||||
server.handle_async_with_channel_join(socket_r, socket_w, user_state.clone(), 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;
|
||||||
if let Ok(user_info) = user_state.user() {
|
if let Ok(user_info) = user_state.user() {
|
||||||
queue.leave_queue(user_info.as_ref().as_ref()).await;
|
queue.leave_queue(user_info.clone()).await;
|
||||||
} else {
|
} else {
|
||||||
log::debug!("Unauthenticated user disconnected");
|
log::debug!("Unauthenticated user disconnected");
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -42,7 +42,7 @@ impl <C: Send + 'static> SimpleOperation<C> for QueueJoinProvider {
|
|||||||
current_lobby.mode,
|
current_lobby.mode,
|
||||||
current_lobby.visibility,
|
current_lobby.visibility,
|
||||||
current_lobby.auto_heal,
|
current_lobby.auto_heal,
|
||||||
user_info.as_ref().as_ref(),
|
user_info.clone(),
|
||||||
events.to_owned(),
|
events.to_owned(),
|
||||||
).await;
|
).await;
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user