Added metrics
This commit is contained in:
+10
-3
@@ -1,5 +1,7 @@
|
||||
/* This file implements all events for the Discord gateway */
|
||||
|
||||
use super::commands::CommandManager;
|
||||
use crate::consts::MOTD;
|
||||
use log::*;
|
||||
use serenity::{
|
||||
async_trait,
|
||||
@@ -13,9 +15,8 @@ use serenity::{
|
||||
prelude::{Context, EventHandler},
|
||||
};
|
||||
|
||||
use crate::consts::MOTD;
|
||||
|
||||
use super::commands::CommandManager;
|
||||
#[cfg(feature = "metrics")]
|
||||
use crate::metrics::MetricsManager;
|
||||
|
||||
// If the GUILD_ID environment variable is set, only allow commands from that guild
|
||||
macro_rules! enforce_guild {
|
||||
@@ -101,6 +102,12 @@ impl Handler {
|
||||
let data = ctx.data.read().await;
|
||||
let command_manager = data.get::<CommandManager>().expect("to contain a value");
|
||||
|
||||
#[cfg(feature = "metrics")]
|
||||
{
|
||||
let metrics = data.get::<MetricsManager>().expect("to contain a value");
|
||||
metrics.command_exec(&command.data.name);
|
||||
}
|
||||
|
||||
command_manager.execute_command(&ctx, command).await;
|
||||
}
|
||||
|
||||
|
||||
+1
-1
@@ -1,5 +1,5 @@
|
||||
pub const VERSION: &str = env!("CARGO_PKG_VERSION");
|
||||
pub const MOTD: &str = "OPEN BETA (v2)";
|
||||
pub const MOTD: &str = "PRE-RELEASE (v2)";
|
||||
|
||||
/// The time it takes for Spoticord to disconnect when no music is being played
|
||||
pub const DISCONNECT_TIME: u64 = 5 * 60;
|
||||
|
||||
+38
-28
@@ -1,18 +1,20 @@
|
||||
use dotenv::dotenv;
|
||||
|
||||
use crate::{bot::commands::CommandManager, database::Database, session::manager::SessionManager};
|
||||
use log::*;
|
||||
use serenity::{framework::StandardFramework, prelude::GatewayIntents, Client};
|
||||
use songbird::SerenityInit;
|
||||
use std::{any::Any, env, process::exit};
|
||||
|
||||
use crate::{
|
||||
bot::commands::CommandManager, database::Database, session::manager::SessionManager,
|
||||
stats::StatsManager,
|
||||
};
|
||||
#[cfg(feature = "metrics")]
|
||||
use metrics::MetricsManager;
|
||||
|
||||
#[cfg(unix)]
|
||||
use tokio::signal::unix::SignalKind;
|
||||
|
||||
#[cfg(feature = "metrics")]
|
||||
mod metrics;
|
||||
|
||||
mod audio;
|
||||
mod bot;
|
||||
mod consts;
|
||||
@@ -21,7 +23,6 @@ mod ipc;
|
||||
mod librespot_ext;
|
||||
mod player;
|
||||
mod session;
|
||||
mod stats;
|
||||
mod utils;
|
||||
|
||||
#[tokio::main]
|
||||
@@ -70,9 +71,13 @@ async fn main() {
|
||||
|
||||
let token = env::var("DISCORD_TOKEN").expect("a token in the environment");
|
||||
let db_url = env::var("DATABASE_URL").expect("a database URL in the environment");
|
||||
let kv_url = env::var("KV_URL").expect("a redis URL in the environment");
|
||||
|
||||
let stats_manager = StatsManager::new(kv_url).expect("Failed to connect to redis");
|
||||
#[cfg(feature = "metrics")]
|
||||
let metrics_manager = {
|
||||
let metrics_url = env::var("METRICS_URL").expect("a prometheus pusher URL in the environment");
|
||||
MetricsManager::new(metrics_url)
|
||||
};
|
||||
|
||||
let session_manager = SessionManager::new();
|
||||
|
||||
// Create client
|
||||
@@ -92,10 +97,13 @@ async fn main() {
|
||||
data.insert::<Database>(Database::new(db_url, None));
|
||||
data.insert::<CommandManager>(CommandManager::new());
|
||||
data.insert::<SessionManager>(session_manager.clone());
|
||||
|
||||
#[cfg(feature = "metrics")]
|
||||
data.insert::<MetricsManager>(metrics_manager.clone());
|
||||
}
|
||||
|
||||
let shard_manager = client.shard_manager.clone();
|
||||
let cache = client.cache_and_http.cache.clone();
|
||||
let _cache = client.cache_and_http.cache.clone();
|
||||
|
||||
#[cfg(unix)]
|
||||
let mut term: Option<Box<dyn Any + Send>> = Some(Box::new(
|
||||
@@ -111,28 +119,27 @@ async fn main() {
|
||||
loop {
|
||||
tokio::select! {
|
||||
_ = tokio::time::sleep(std::time::Duration::from_secs(60)) => {
|
||||
let guild_count = cache.guilds().len();
|
||||
let active_count = session_manager.get_active_session_count().await;
|
||||
let total_count = session_manager.get_session_count().await;
|
||||
#[cfg(feature = "metrics")]
|
||||
{
|
||||
let guild_count = _cache.guilds().len();
|
||||
let active_count = session_manager.get_active_session_count().await;
|
||||
let total_count = session_manager.get_session_count().await;
|
||||
|
||||
if let Err(why) = stats_manager.set_server_count(guild_count) {
|
||||
error!("Failed to update server count: {}", why);
|
||||
metrics_manager.set_server_count(guild_count);
|
||||
metrics_manager.set_active_sessions(active_count);
|
||||
metrics_manager.set_total_sessions(total_count);
|
||||
|
||||
// Yes, I like to handle my s's when I'm working with amounts
|
||||
debug!(
|
||||
"Updated metrics: {} guild{}, {} active session{}, {} total session{}",
|
||||
guild_count,
|
||||
if guild_count == 1 { "" } else { "s" },
|
||||
active_count,
|
||||
if active_count == 1 { "" } else { "s" },
|
||||
total_count,
|
||||
if total_count == 1 { "" } else { "s" }
|
||||
);
|
||||
}
|
||||
|
||||
if let Err(why) = stats_manager.set_active_count(active_count) {
|
||||
error!("Failed to update active count: {}", why);
|
||||
}
|
||||
|
||||
// Yes, I like to handle my s's when I'm working with amounts
|
||||
debug!(
|
||||
"Updated stats: {} guild{}, {} active session{}, {} total session{}",
|
||||
guild_count,
|
||||
if guild_count == 1 { "" } else { "s" },
|
||||
active_count,
|
||||
if active_count == 1 { "" } else { "s" },
|
||||
total_count,
|
||||
if total_count == 1 { "" } else { "s" }
|
||||
);
|
||||
}
|
||||
|
||||
_ = tokio::signal::ctrl_c() => {
|
||||
@@ -159,6 +166,9 @@ async fn main() {
|
||||
|
||||
shard_manager.lock().await.shutdown_all().await;
|
||||
|
||||
#[cfg(feature = "metrics")]
|
||||
metrics_manager.stop();
|
||||
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
+124
@@ -0,0 +1,124 @@
|
||||
use std::{
|
||||
collections::hash_map::RandomState,
|
||||
sync::{
|
||||
atomic::{AtomicBool, Ordering},
|
||||
Arc,
|
||||
},
|
||||
thread,
|
||||
time::Duration,
|
||||
};
|
||||
|
||||
use lazy_static::lazy_static;
|
||||
use prometheus::{
|
||||
opts, push_metrics, register_int_counter_vec, register_int_gauge, IntCounterVec, IntGauge,
|
||||
};
|
||||
use serenity::prelude::TypeMapKey;
|
||||
|
||||
use crate::session::pbi::PlaybackInfo;
|
||||
|
||||
lazy_static! {
|
||||
static ref TOTAL_SERVERS: IntGauge =
|
||||
register_int_gauge!("total_servers", "Total number of servers Spoticord is in").unwrap();
|
||||
static ref ACTIVE_SESSIONS: IntGauge = register_int_gauge!(
|
||||
"active_sessions",
|
||||
"Total number of servers with an active Spoticord session"
|
||||
)
|
||||
.unwrap();
|
||||
static ref TOTAL_SESSIONS: IntGauge = register_int_gauge!(
|
||||
"total_sessions",
|
||||
"Total number of servers with Spoticord in a voice channel"
|
||||
)
|
||||
.unwrap();
|
||||
static ref TRACKS_PLAYED: IntCounterVec = register_int_counter_vec!(
|
||||
opts!("tracks_played", "Tracks Played"),
|
||||
&["type", "name", "artists", "uri"]
|
||||
)
|
||||
.unwrap();
|
||||
static ref COMMANDS_EXECUTED: IntCounterVec = register_int_counter_vec!(
|
||||
opts!("commands_executed", "Commands Executed"),
|
||||
&["command"]
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct MetricsManager {
|
||||
should_stop: Arc<AtomicBool>,
|
||||
}
|
||||
|
||||
impl MetricsManager {
|
||||
pub fn new(pusher_url: impl Into<String>) -> Self {
|
||||
let instance = Self {
|
||||
should_stop: Arc::new(AtomicBool::new(false)),
|
||||
};
|
||||
|
||||
thread::spawn({
|
||||
let instance = instance.clone();
|
||||
let pusher_url = pusher_url.into();
|
||||
|
||||
move || loop {
|
||||
thread::sleep(Duration::from_secs(5));
|
||||
|
||||
if instance.should_stop() {
|
||||
break;
|
||||
}
|
||||
|
||||
if let Err(why) = push_metrics::<RandomState>(
|
||||
"spoticord_metrics",
|
||||
Default::default(),
|
||||
&pusher_url,
|
||||
prometheus::gather(),
|
||||
None,
|
||||
) {
|
||||
log::error!("Failed to push metrics: {}", why);
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
instance
|
||||
}
|
||||
|
||||
pub fn should_stop(&self) -> bool {
|
||||
self.should_stop.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
pub fn stop(&self) {
|
||||
self.should_stop.store(true, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
pub fn set_server_count(&self, count: usize) {
|
||||
TOTAL_SERVERS.set(count as i64);
|
||||
}
|
||||
|
||||
pub fn set_total_sessions(&self, count: usize) {
|
||||
TOTAL_SESSIONS.set(count as i64);
|
||||
}
|
||||
|
||||
pub fn set_active_sessions(&self, count: usize) {
|
||||
ACTIVE_SESSIONS.set(count as i64);
|
||||
}
|
||||
|
||||
pub fn track_play(&self, track: &PlaybackInfo) {
|
||||
let track_type = match track.get_type() {
|
||||
Some(track_type) => track_type,
|
||||
None => return,
|
||||
};
|
||||
|
||||
TRACKS_PLAYED
|
||||
.with_label_values(&[
|
||||
&track_type,
|
||||
&track.get_name().expect("To have a name"),
|
||||
&track.get_artists().expect("To have artists"),
|
||||
&track.get_url().expect("To have a URL"),
|
||||
])
|
||||
.inc();
|
||||
}
|
||||
|
||||
pub fn command_exec(&self, command: &str) {
|
||||
COMMANDS_EXECUTED.with_label_values(&[command]).inc();
|
||||
}
|
||||
}
|
||||
|
||||
impl TypeMapKey for MetricsManager {
|
||||
type Value = MetricsManager;
|
||||
}
|
||||
@@ -173,11 +173,13 @@ impl SessionManager {
|
||||
}
|
||||
|
||||
/// Get the amount of sessions
|
||||
#[allow(dead_code)]
|
||||
pub async fn get_session_count(&self) -> usize {
|
||||
self.0.read().await.get_session_count()
|
||||
}
|
||||
|
||||
/// Get the amount of sessions with an owner
|
||||
#[allow(dead_code)]
|
||||
pub async fn get_active_session_count(&self) -> usize {
|
||||
self.0.read().await.get_active_session_count().await
|
||||
}
|
||||
|
||||
@@ -33,6 +33,9 @@ use std::{
|
||||
};
|
||||
use tokio::sync::Mutex;
|
||||
|
||||
#[cfg(feature = "metrics")]
|
||||
use crate::metrics::MetricsManager;
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct SpoticordSession(Arc<RwLock<InnerSpoticordSession>>);
|
||||
|
||||
@@ -58,6 +61,9 @@ struct InnerSpoticordSession {
|
||||
/// Whether the session has been disconnected
|
||||
/// If this is true then this instance should no longer be used and dropped
|
||||
disconnected: bool,
|
||||
|
||||
#[cfg(feature = "metrics")]
|
||||
metrics: MetricsManager,
|
||||
}
|
||||
|
||||
impl SpoticordSession {
|
||||
@@ -75,6 +81,12 @@ impl SpoticordSession {
|
||||
.expect("to contain a value")
|
||||
.clone();
|
||||
|
||||
#[cfg(feature = "metrics")]
|
||||
let metrics = data
|
||||
.get::<MetricsManager>()
|
||||
.expect("to contain a value")
|
||||
.clone();
|
||||
|
||||
// Join the voice channel
|
||||
let songbird = songbird::get(ctx).await.expect("to be present").clone();
|
||||
|
||||
@@ -98,6 +110,9 @@ impl SpoticordSession {
|
||||
disconnect_handle: None,
|
||||
client: None,
|
||||
disconnected: false,
|
||||
|
||||
#[cfg(feature = "metrics")]
|
||||
metrics,
|
||||
};
|
||||
|
||||
let mut instance = Self(Arc::new(RwLock::new(inner)));
|
||||
@@ -522,6 +537,14 @@ impl SpoticordSession {
|
||||
pbi.update_track_episode(spotify_id, track, episode);
|
||||
}
|
||||
|
||||
// Send track play event to metrics
|
||||
#[cfg(feature = "metrics")]
|
||||
{
|
||||
if let Some(ref pbi) = inner.playback_info {
|
||||
inner.metrics.track_play(pbi);
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
@@ -106,4 +106,28 @@ impl PlaybackInfo {
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
/// Get the type of audio (track or episode)
|
||||
#[allow(dead_code)]
|
||||
pub fn get_type(&self) -> Option<String> {
|
||||
if self.track.is_some() {
|
||||
Some("track".into())
|
||||
} else if self.episode.is_some() {
|
||||
Some("episode".into())
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
/// Get the public facing url of the track or episode
|
||||
#[allow(dead_code)]
|
||||
pub fn get_url(&self) -> Option<&str> {
|
||||
if let Some(ref track) = self.track {
|
||||
Some(track.external_urls.spotify.as_str())
|
||||
} else if let Some(ref episode) = self.episode {
|
||||
Some(episode.external_urls.spotify.as_str())
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,26 +0,0 @@
|
||||
use redis::{Commands, RedisResult};
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct StatsManager {
|
||||
redis: redis::Client,
|
||||
}
|
||||
|
||||
impl StatsManager {
|
||||
pub fn new(url: impl Into<String>) -> RedisResult<StatsManager> {
|
||||
let redis = redis::Client::open(url.into())?;
|
||||
|
||||
Ok(StatsManager { redis })
|
||||
}
|
||||
|
||||
pub fn set_server_count(&self, count: usize) -> RedisResult<()> {
|
||||
let mut con = self.redis.get_connection()?;
|
||||
|
||||
con.set("sc-bot-total-servers", count.to_string())
|
||||
}
|
||||
|
||||
pub fn set_active_count(&self, count: usize) -> RedisResult<()> {
|
||||
let mut con = self.redis.get_connection()?;
|
||||
|
||||
con.set("sc-bot-active-servers", count.to_string())
|
||||
}
|
||||
}
|
||||
@@ -23,11 +23,17 @@ pub struct Album {
|
||||
pub images: Vec<Image>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Deserialize)]
|
||||
pub struct ExternalUrls {
|
||||
pub spotify: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Deserialize)]
|
||||
pub struct Track {
|
||||
pub name: String,
|
||||
pub artists: Vec<Artist>,
|
||||
pub album: Album,
|
||||
pub external_urls: ExternalUrls,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Deserialize)]
|
||||
@@ -40,6 +46,7 @@ pub struct Show {
|
||||
pub struct Episode {
|
||||
pub name: String,
|
||||
pub show: Show,
|
||||
pub external_urls: ExternalUrls,
|
||||
}
|
||||
|
||||
pub async fn get_username(token: impl Into<String>) -> Result<String, String> {
|
||||
|
||||
Reference in New Issue
Block a user