mirror of
https://git.ngram.ca/OpenJam/rc-servers
synced 2026-08-23 23:08:52 +00:00
Track and report lobby queue wait times
This commit is contained in:
@@ -128,6 +128,7 @@ pub struct QueueHandler {
|
||||
autostart_after: Option<std::time::Duration>,
|
||||
autostart_task_started: std::sync::atomic::AtomicBool,
|
||||
team_choosers: std::sync::Arc<crate::team_selection::InitedTeamChoosers>,
|
||||
wait_times: std::sync::Arc<super::queue_time_tracker::QueueTimeTracker>,
|
||||
}
|
||||
|
||||
impl QueueHandler {
|
||||
@@ -167,6 +168,7 @@ impl QueueHandler {
|
||||
autostart_after: mp_settings.lobby_autostart_after,
|
||||
autostart_task_started: std::sync::atomic::AtomicBool::new(false),
|
||||
team_choosers: std::sync::Arc::new(team_choosers),
|
||||
wait_times: std::sync::Arc::new(super::queue_time_tracker::QueueTimeTracker::new()),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -185,6 +187,7 @@ impl QueueHandler {
|
||||
let weapon_guesser = self.weapon_guesser.clone();
|
||||
let autostart_after = self.autostart_after.unwrap();
|
||||
let team_choosers = self.team_choosers.clone();
|
||||
let wait_times = self.wait_times.clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
loop {
|
||||
@@ -250,6 +253,7 @@ impl QueueHandler {
|
||||
q_entry,
|
||||
starter.as_ref().as_ref(),
|
||||
team_choosers.as_ref(),
|
||||
wait_times.as_ref(),
|
||||
).await;
|
||||
}
|
||||
|
||||
@@ -280,7 +284,8 @@ impl QueueHandler {
|
||||
key,
|
||||
q_entry,
|
||||
user,
|
||||
&self.team_choosers
|
||||
&self.team_choosers,
|
||||
&self.wait_times,
|
||||
).await
|
||||
}
|
||||
|
||||
@@ -310,6 +315,7 @@ impl QueueHandler {
|
||||
mut q_entry: Queue,
|
||||
user: &(dyn oj_rc_core::persist::user::LobbyUser + Send + Sync),
|
||||
team_choosers: &crate::team_selection::InitedTeamChoosers,
|
||||
wait_times: &crate::queue_time_tracker::QueueTimeTracker,
|
||||
) {
|
||||
let guid_str = key.unique_guid();
|
||||
let game_desc = oj_rc_core::persist::user::GameDescriptor {
|
||||
@@ -339,12 +345,14 @@ impl QueueHandler {
|
||||
player_descs.push(lobby_desc);
|
||||
}
|
||||
|
||||
wait_times.update_time_match_starting_now(q_entry.users.iter().map(|x| x.enqueued_at.clone()));
|
||||
|
||||
let missing = users_per_game.saturating_sub(q_entry.users.len());
|
||||
|
||||
match user.start_game(game_desc, player_descs, factory.as_ref(), &cpu_counter, &weapon_guesser, team_picker, missing).await {
|
||||
Ok(fakes) => {
|
||||
let player_datas = q_entry.users.iter().map(|x| x.player.clone())
|
||||
.chain(fakes.players.into_iter().map(|(desc, _emu)| desc),)
|
||||
.chain(fakes.players.into_iter().map(|(desc, _emu)| desc))
|
||||
.collect();
|
||||
let enter_battle_ev = crate::events::battle_enter::BattleEnter {
|
||||
host: hostname.to_string(),
|
||||
@@ -805,4 +813,8 @@ impl QueueHandler {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub fn wait_time_s(&self) -> i32 {
|
||||
self.wait_times.get_average()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,6 +7,7 @@ mod data;
|
||||
mod operations;
|
||||
mod events;
|
||||
mod team_selection;
|
||||
mod queue_time_tracker;
|
||||
|
||||
use oj_polariton_auth::Handshake;
|
||||
use tokio::net;
|
||||
|
||||
@@ -37,7 +37,7 @@ impl <C: Send + 'static> SimpleOperation<C> for QueueJoinProvider {
|
||||
},
|
||||
oj_rc_core::data::lobby::LobbyType::CustomGame => {
|
||||
log::debug!("Joining platoon {} to custom game lobby", group_id.string);
|
||||
params.insert(ESTIMATED_QUEUE_TIME_PARAM_KEY, Typed::Int(42));
|
||||
params.insert(ESTIMATED_QUEUE_TIME_PARAM_KEY, Typed::Int(self.queue_handler.wait_time_s()));
|
||||
params.insert(PERSONAL_RANKING_PARAM_KEY, Typed::Double(42.0));
|
||||
let events = user.event_sender();
|
||||
let user_info = user.user()?;
|
||||
@@ -50,7 +50,7 @@ impl <C: Send + 'static> SimpleOperation<C> for QueueJoinProvider {
|
||||
// regular multiplayer
|
||||
if let Some(Typed::Str(event_to_join)) = params.remove(&EVENT_TO_JOIN_PARAM_KEY) {
|
||||
log::debug!("Joining platoon {} to multiplayer lobby with event {}", group_id.string, event_to_join.string);
|
||||
params.insert(ESTIMATED_QUEUE_TIME_PARAM_KEY, Typed::Int(42));
|
||||
params.insert(ESTIMATED_QUEUE_TIME_PARAM_KEY, Typed::Int(self.queue_handler.wait_time_s()));
|
||||
params.insert(PERSONAL_RANKING_PARAM_KEY, Typed::Double(42.0));
|
||||
let events = user.event_sender();
|
||||
let user_info = user.user()?;
|
||||
|
||||
57
rc_lobby_room/src/queue_time_tracker.rs
Normal file
57
rc_lobby_room/src/queue_time_tracker.rs
Normal file
@@ -0,0 +1,57 @@
|
||||
const MAX_HISTORY_LEN: usize = 64;
|
||||
|
||||
struct QueueTime {
|
||||
time: std::time::Duration,
|
||||
players: usize,
|
||||
}
|
||||
|
||||
pub struct QueueTimeTracker {
|
||||
average: std::sync::atomic::AtomicI32, // seconds
|
||||
history: std::sync::Mutex<std::collections::VecDeque<QueueTime>>,
|
||||
}
|
||||
|
||||
impl QueueTimeTracker {
|
||||
pub fn new() -> Self {
|
||||
Self {
|
||||
average: std::sync::atomic::AtomicI32::new(42),
|
||||
history: std::sync::Mutex::new(std::collections::VecDeque::with_capacity(MAX_HISTORY_LEN)),
|
||||
}
|
||||
}
|
||||
|
||||
fn update_time(&self, players_enqueue_at: impl Iterator<Item = chrono::DateTime<chrono::Utc>>, game_start: chrono::DateTime<chrono::Utc>) {
|
||||
let mut total_duration = std::time::Duration::ZERO;
|
||||
let mut total_players: usize = 0;
|
||||
for player_enqueue_at in players_enqueue_at {
|
||||
total_players += 1;
|
||||
total_duration += game_start.signed_duration_since(player_enqueue_at).to_std().unwrap_or(std::time::Duration::ZERO);
|
||||
}
|
||||
let time = total_duration.div_f64(total_players as f64);
|
||||
let queue_time = QueueTime {
|
||||
time,
|
||||
players: total_players,
|
||||
};
|
||||
let mut lock = self.history.lock().unwrap();
|
||||
if lock.len() == MAX_HISTORY_LEN {
|
||||
lock.pop_front();
|
||||
}
|
||||
lock.push_back(queue_time);
|
||||
let mut total_time: u128 = 0;
|
||||
let mut total_players: u128 = 0;
|
||||
for time in lock.iter() {
|
||||
total_players += time.players as u128;
|
||||
total_time += (time.time.as_secs() as u128) * (time.players as u128);
|
||||
}
|
||||
let average = total_time / total_players;
|
||||
let old_average = self.average.swap(average.try_into().unwrap_or_default(), std::sync::atomic::Ordering::Relaxed);
|
||||
log::debug!("Average queue time is now {}s, was {}s", average, old_average);
|
||||
}
|
||||
|
||||
pub fn update_time_match_starting_now(&self, players_enqueue_at: impl Iterator<Item = chrono::DateTime<chrono::Utc>>) {
|
||||
let now = chrono::Utc::now();
|
||||
self.update_time(players_enqueue_at, now);
|
||||
}
|
||||
|
||||
pub fn get_average(&self) -> i32 {
|
||||
self.average.load(std::sync::atomic::Ordering::Relaxed)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user