Initial v2.2.0 push

This commit is contained in:
DaXcess
2024-08-12 16:14:24 +02:00
parent d1c01e9274
commit 97c18385dd
97 changed files with 7252 additions and 5211 deletions
+21
View File
@@ -0,0 +1,21 @@
[package]
name = "spoticord_session"
version = "2.2.0"
edition = "2021"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
[dependencies]
spoticord_config = { path = "../spoticord_config" }
spoticord_database = { path = "../spoticord_database" }
spoticord_player = { path = "../spoticord_player" }
spoticord_utils = { path = "../spoticord_utils" }
tokio = { version = "1.39.2", features = ["full"] }
librespot = { git = "https://github.com/SpoticordMusic/librespot.git", version = "0.5.0-dev", default-features = false }
serenity = "0.12.2"
songbird = { version = "0.4.3", features = ["simd-json"] }
anyhow = "1.0.86"
log = "0.4.22"
base64 = "0.22.1"
poise = "0.6.1"
+580
View File
@@ -0,0 +1,580 @@
pub mod lyrics_embed;
pub mod manager;
pub mod playback_embed;
use anyhow::{anyhow, Result};
use base64::{engine::general_purpose::STANDARD as BASE64, Engine};
use librespot::{discovery::Credentials, protocol::authentication::AuthenticationType};
use log::{debug, error, trace};
use lyrics_embed::LyricsEmbed;
use manager::{SessionManager, SessionQuery};
use playback_embed::{PlaybackEmbed, PlaybackEmbedHandle};
use serenity::{
all::{
ChannelId, CommandInteraction, CreateEmbed, CreateMessage, GuildChannel, GuildId, UserId,
},
async_trait,
};
use songbird::{model::payload::ClientDisconnect, Call, CoreEvent, Event, EventContext};
use spoticord_database::Database;
use spoticord_player::{Player, PlayerEvent, PlayerHandle};
use spoticord_utils::{discord::Colors, spotify};
use std::{ops::ControlFlow, sync::Arc, time::Duration};
use tokio::{
sync::{mpsc, oneshot, Mutex},
task::JoinHandle,
};
#[derive(Debug)]
pub enum SessionCommand {
GetOwner(oneshot::Sender<UserId>),
GetPlayer(oneshot::Sender<PlayerHandle>),
GetActive(oneshot::Sender<bool>),
CreatePlaybackEmbed(SessionHandle, CommandInteraction),
CreateLyricsEmbed(SessionHandle, CommandInteraction),
Reactivate(UserId, oneshot::Sender<Result<()>>),
ShutdownPlayer,
Disconnect,
DisconnectTimedOut,
}
pub struct Session {
session_manager: SessionManager,
context: serenity::all::Context,
guild_id: GuildId,
text_channel: GuildChannel,
call: Arc<Mutex<Call>>,
player: PlayerHandle,
owner: UserId,
active: bool,
timeout_tx: Option<oneshot::Sender<()>>,
commands: mpsc::Receiver<SessionCommand>,
events: mpsc::Receiver<PlayerEvent>,
commands_inner_tx: mpsc::Sender<SessionCommand>,
commands_inner_rx: mpsc::Receiver<SessionCommand>,
playback_embed: Option<PlaybackEmbedHandle>,
lyrics_embed: Option<JoinHandle<()>>,
}
impl Session {
pub async fn create(
session_manager: SessionManager,
context: &serenity::all::Context,
guild_id: GuildId,
voice_channel_id: ChannelId,
text_channel_id: ChannelId,
owner: UserId,
) -> Result<SessionHandle> {
// Set up communication channel
let (tx, rx) = mpsc::channel(16);
let handle = SessionHandle {
guild: guild_id,
voice_channel: voice_channel_id,
text_channel: text_channel_id,
commands: tx,
};
// Resolve text channel
let text_channel = text_channel_id
.to_channel(&context)
.await?
.guild()
.ok_or(anyhow!("Text channel is not a guild channel"))?;
// Create channel for internal command communication (timeouts hint hint)
// This uses separate channels as to not cause a cyclic dependency
let (inner_tx, inner_rx) = mpsc::channel(16);
// Grab user credentials and info before joining call
let credentials =
retrieve_credentials(&session_manager.database(), owner.to_string()).await?;
let device_name = session_manager
.database()
.get_user(owner.to_string())
.await?
.device_name;
// Hello Discord I'm here
let call = session_manager
.songbird()
.join(guild_id, voice_channel_id)
.await?;
// Make sure call guard is dropped or else we can't execute session.run
{
let mut call = call.lock().await;
// Wasn't able to confirm if this is true, but this might reduce network bandwith by not receiving user voice packets
_ = call.deafen(true).await;
// Set up call events
call.add_global_event(Event::Core(CoreEvent::DriverDisconnect), handle.clone());
call.add_global_event(Event::Core(CoreEvent::ClientDisconnect), handle.clone());
}
let (player, events) = match Player::create(credentials, call.clone(), device_name).await {
Ok(player) => player,
Err(why) => {
// Leave call on error, otherwise bot will be stuck in call forever until manually disconnected or taken over
_ = call.lock().await.leave().await;
error!("Failed to create player: {why}");
return Err(why);
}
};
let mut session = Self {
session_manager,
context: context.to_owned(),
text_channel,
call,
player,
guild_id,
owner,
active: true,
timeout_tx: None,
commands: rx,
events,
commands_inner_tx: inner_tx,
commands_inner_rx: inner_rx,
playback_embed: None,
lyrics_embed: None,
};
session.start_timeout();
tokio::spawn(session.run());
Ok(handle)
}
pub async fn run(mut self) {
loop {
tokio::select! {
opt_command = self.commands.recv() => {
let Some(command) = opt_command else {
break;
};
if self.handle_command(command).await.is_break() {
break;
}
},
opt_event = self.events.recv(), if self.active => {
let Some(event) = opt_event else {
self.shutdown_player().await;
continue;
};
self.handle_event(event).await;
},
// Internal communication channel
Some(command) = self.commands_inner_rx.recv() => {
if self.handle_command(command).await.is_break() {
break;
}
}
else => break,
}
}
}
async fn handle_command(&mut self, command: SessionCommand) -> ControlFlow<(), ()> {
trace!("SessionCommand::{command:?}");
match command {
SessionCommand::GetOwner(sender) => _ = sender.send(self.owner),
SessionCommand::GetPlayer(sender) => _ = sender.send(self.player.clone()),
SessionCommand::GetActive(sender) => _ = sender.send(self.active),
SessionCommand::CreatePlaybackEmbed(handle, interaction) => {
match PlaybackEmbed::create(self, handle, interaction).await {
Ok(Some(playback_embed)) => {
self.playback_embed = Some(playback_embed);
}
Ok(None) => {}
Err(why) => {
error!("Failed to create playing embed: {why}");
}
};
}
SessionCommand::CreateLyricsEmbed(handle, interaction) => {
match LyricsEmbed::create(self, handle, interaction).await {
Ok(Some(lyrics_embed)) => {
if let Some(current) = self.lyrics_embed.take() {
current.abort();
}
self.lyrics_embed = Some(lyrics_embed);
}
Ok(None) => {}
Err(why) => {
error!("Failed to create lyrics embed: {why}");
}
}
}
SessionCommand::Reactivate(new_owner, tx) => {
_ = tx.send(self.reactivate(new_owner).await)
}
SessionCommand::ShutdownPlayer => self.shutdown_player().await,
SessionCommand::Disconnect => {
self.disconnect().await;
return ControlFlow::Break(());
}
SessionCommand::DisconnectTimedOut => {
self.disconnect().await;
_ = self
.text_channel
.send_message(
&self.context,
CreateMessage::new().embed(
CreateEmbed::new()
.title("It's a little quiet in here")
.description("The bot has been inactive for too long, and has been disconnected.")
.color(Colors::Warning),
),
)
.await;
return ControlFlow::Break(());
}
};
ControlFlow::Continue(())
}
async fn handle_event(&mut self, event: PlayerEvent) {
match event {
PlayerEvent::Play => self.stop_timeout(),
PlayerEvent::Pause => self.start_timeout(),
PlayerEvent::Stopped => self.shutdown_player().await,
PlayerEvent::TrackChanged(_) => {}
}
if let Some(playback_embed) = &self.playback_embed {
if playback_embed.invoke_update().await.is_err() {
self.playback_embed = None;
}
}
}
fn start_timeout(&mut self) {
if let Some(tx) = self.timeout_tx.take() {
_ = tx.send(());
}
let (tx, rx) = oneshot::channel::<()>();
self.timeout_tx = Some(tx);
let inner_tx = self.commands_inner_tx.clone();
tokio::spawn(async move {
let mut timer =
tokio::time::interval(Duration::from_secs(spoticord_config::DISCONNECT_TIME));
// Ignore immediate tick
timer.tick().await;
tokio::select! {
_ = rx => return,
_ = timer.tick() => {}
};
// Disconnect through inner communication
_ = inner_tx.send(SessionCommand::DisconnectTimedOut).await;
});
}
fn stop_timeout(&mut self) {
if let Some(tx) = self.timeout_tx.take() {
_ = tx.send(());
}
}
async fn reactivate(&mut self, new_owner: UserId) -> Result<()> {
if self.active {
return Err(anyhow!("Cannot reactivate session that is already active"));
}
let credentials =
retrieve_credentials(&self.session_manager.database(), new_owner.to_string()).await?;
let device_name = self
.session_manager
.database()
.get_user(new_owner.to_string())
.await?
.device_name;
let (player, player_events) =
Player::create(credentials, self.call.clone(), device_name).await?;
self.owner = new_owner;
self.player = player;
self.events = player_events;
self.active = true;
Ok(())
}
async fn shutdown_player(&mut self) {
self.player.shutdown().await;
self.start_timeout();
self.active = false;
// Remove owner from session manager
self.session_manager
.remove_session(SessionQuery::Owner(self.owner));
}
async fn disconnect(&mut self) {
// Kill timeout if one is running
self.stop_timeout();
// Force close channels, as handles may otherwise hold this struct hostage
self.commands.close();
self.events.close();
// Leave call, ignore errors
let mut call = self.call.lock().await;
_ = call.leave().await;
}
}
impl Drop for Session {
fn drop(&mut self) {
// Abort timeout task
if let Some(tx) = self.timeout_tx.take() {
_ = tx.send(());
}
// Abort lyrics task
if let Some(lyrics) = self.lyrics_embed.take() {
lyrics.abort();
}
// Clean up the session from the session manager
// This is done in Drop::drop to ensure that the session always cleans up after itself
// even if something went wrong
let session_manager = self.session_manager.clone();
let guild_id = self.guild_id;
let owner = self.owner;
session_manager.remove_session(SessionQuery::Guild(guild_id));
session_manager.remove_session(SessionQuery::Owner(owner));
}
}
#[derive(Clone, Debug)]
pub struct SessionHandle {
guild: GuildId,
voice_channel: ChannelId,
text_channel: ChannelId,
commands: mpsc::Sender<SessionCommand>,
}
impl SessionHandle {
/// Check if the session handle is valid
pub fn is_valid(&self) -> bool {
!self.commands.is_closed()
}
pub fn guild(&self) -> GuildId {
self.guild
}
pub fn voice_channel(&self) -> ChannelId {
self.voice_channel
}
pub fn text_channel(&self) -> ChannelId {
self.text_channel
}
/// Retrieve the current owner of the session
pub async fn owner(&self) -> Result<UserId> {
let (tx, rx) = oneshot::channel();
self.commands.send(SessionCommand::GetOwner(tx)).await?;
let result = rx.await?;
Ok(result)
}
/// Retrieve the player handle from the session
pub async fn player(&self) -> Result<PlayerHandle> {
let (tx, rx) = oneshot::channel();
self.commands.send(SessionCommand::GetPlayer(tx)).await?;
let result = rx.await?;
Ok(result)
}
pub async fn active(&self) -> Result<bool> {
let (tx, rx) = oneshot::channel();
self.commands.send(SessionCommand::GetActive(tx)).await?;
let result = rx.await?;
Ok(result)
}
/// Instruct the session to make another user owner.
///
/// This will fail if the session still has an active user assigned to it.
pub async fn reactivate(&self, new_owner: UserId) -> Result<()> {
let (tx, rx) = oneshot::channel();
self.commands
.send(SessionCommand::Reactivate(new_owner, tx))
.await?;
rx.await?
}
/// Create a playback embed as a response to an interaction
///
/// This playback embed will automatically update when certain events happen
pub async fn create_playback_embed(&self, interaction: CommandInteraction) -> Result<()> {
self.commands
.send(SessionCommand::CreatePlaybackEmbed(
self.clone(),
interaction,
))
.await?;
Ok(())
}
/// Create a lyrics embed as a response to an interaction
///
/// This lyrics embed will automatically retrieve the lyrics and update the embed accordingly
pub async fn create_lyrics_embed(&self, interaction: CommandInteraction) -> Result<()> {
self.commands
.send(SessionCommand::CreateLyricsEmbed(self.clone(), interaction))
.await?;
Ok(())
}
/// Instruct the session to destroy the player (but keep voice call).
///
/// This is meant to be used for when the session owner leaves the call
/// and allows other users to become owner using the `/join` command.
///
/// This should also remove the owner from the session manager.
pub async fn shutdown_player(&self) {
if let Err(why) = self.commands.send(SessionCommand::ShutdownPlayer).await {
error!("Failed to send command: {why}");
}
}
/// Instruct the session to destroy itself.
///
/// This should also remove the player and the owner from the session manager.
pub async fn disconnect(&self) {
if let Err(why) = self.commands.send(SessionCommand::Disconnect).await {
error!("Failed to send command: {why}");
}
}
}
#[async_trait]
impl songbird::EventHandler for SessionHandle {
async fn act(&self, event: &EventContext<'_>) -> Option<Event> {
if !self.is_valid() {
return Some(Event::Cancel);
}
match event {
EventContext::DriverDisconnect(_) => {
debug!("Bot disconnected from call, cleaning up");
self.disconnect().await;
}
EventContext::ClientDisconnect(ClientDisconnect { user_id }) => {
match self.owner().await {
Ok(id) if id.get() == user_id.0 => {
debug!("Owner of session disconnected, stopping playback");
self.shutdown_player().await;
}
_ => {}
}
}
_ => {}
}
None
}
}
async fn retrieve_credentials(database: &Database, owner: impl AsRef<str>) -> Result<Credentials> {
let account = database.get_account(&owner).await?;
let token = if let Some(session_token) = &account.session_token {
match spotify::validate_token(&account.username, session_token).await {
Ok(Some(token)) => {
database
.update_session_token(&account.user_id, &token)
.await?;
Some(token)
}
Ok(None) => Some(session_token.clone()),
Err(_) => None,
}
} else {
None
};
// Request new session token if previous one was invalid or missing
let token = match token {
Some(token) => token,
None => {
let access_token = database.get_access_token(&account.user_id).await?;
let credentials = spotify::request_session_token(Credentials {
username: account.username.clone(),
auth_type: AuthenticationType::AUTHENTICATION_SPOTIFY_TOKEN,
auth_data: access_token.into_bytes(),
})
.await?;
let token = BASE64.encode(credentials.auth_data);
database
.update_session_token(&account.user_id, &token)
.await?;
token
}
};
Ok(Credentials {
username: account.username,
auth_type: AuthenticationType::AUTHENTICATION_STORED_SPOTIFY_CREDENTIALS,
auth_data: BASE64.decode(token)?,
})
}
+447
View File
@@ -0,0 +1,447 @@
use std::{ops::ControlFlow, time::Duration};
use anyhow::Result;
use librespot::{
core::SpotifyId,
metadata::{
lyrics::{Line, SyncType},
Lyrics,
},
};
use log::error;
use serenity::{
all::{
CommandInteraction, ComponentInteraction, ComponentInteractionCollector, Context,
CreateActionRow, CreateButton, CreateEmbed, CreateEmbedFooter, CreateInteractionResponse,
CreateInteractionResponseMessage, EditMessage, Message,
},
futures::StreamExt,
};
use spoticord_player::info::PlaybackInfo;
use spoticord_utils::discord::Colors;
use tokio::task::JoinHandle;
use crate::{Session, SessionHandle};
const PAGE_LENGTH: usize = 3000;
const TIME_OFFSET: u32 = 1000;
pub struct LyricsEmbed {
guild_id: String,
ctx: Context,
session: SessionHandle,
message: Message,
track: SpotifyId,
lyrics: Option<Lyrics>,
page: usize,
}
impl LyricsEmbed {
pub async fn create(
session: &Session,
handle: SessionHandle,
interaction: CommandInteraction,
) -> Result<Option<JoinHandle<()>>> {
let ctx = session.context.clone();
if !session.active {
respond_not_playing(&ctx, interaction).await?;
return Ok(None);
}
let Some(playback_info) = session.player.playback_info().await? else {
respond_not_playing(&ctx, interaction).await?;
return Ok(None);
};
let guild_id = interaction
.guild_id
.expect("interaction was outside of a guild")
.to_string();
let lyrics = session.player.get_lyrics().await?;
// Send initial message
interaction
.create_response(
&ctx,
CreateInteractionResponse::Message(
CreateInteractionResponseMessage::new()
.embed(lyrics_embed(&lyrics, &playback_info, 0))
.components(vec![lyrics_buttons(&guild_id, &lyrics, 0)]),
),
)
.await?;
// Retrieve message instead of editing interaction response, as those tokens are only valid for 15 minutes
let message = interaction.get_response(&ctx).await?;
let this = Self {
guild_id: guild_id.clone(),
ctx: ctx.clone(),
session: handle,
message,
track: playback_info.track_id(),
lyrics,
page: 0,
};
let collector = ComponentInteractionCollector::new(&ctx)
.filter(move |press| {
let parts = press.data.custom_id.split(':').collect::<Vec<_>>();
matches!(parts.first(), Some(&"lyrics"))
&& matches!(parts.last(), Some(id) if id == &guild_id)
})
.timeout(Duration::from_secs(3600 * 24));
let handle = tokio::spawn(this.run(collector));
Ok(Some(handle))
}
async fn run(mut self, collector: ComponentInteractionCollector) {
let mut stream = collector.stream();
let mut interval = tokio::time::interval(Duration::from_secs(2));
loop {
tokio::select! {
_ = interval.tick() => {
if self.handle_tick().await.is_break() {
break;
}
}
opt_press = stream.next() => {
let Some(press) = opt_press else {
break;
};
// Immediately acknowledge, we don't have to inform the user about the update
_ = press
.create_response(&self.ctx, CreateInteractionResponse::Acknowledge)
.await;
if self.handle_press(press).await.is_break() {
break;
}
}
}
}
}
async fn handle_tick(&mut self) -> ControlFlow<(), ()> {
let Ok(player) = self.session.player().await else {
// Failure means that the session is gone, so we quit
return ControlFlow::Break(());
};
if !matches!(self.session.active().await, Ok(true)) {
// If the session is currently not active, just wait until it becomes active again
return ControlFlow::Continue(());
}
let Ok(Some(playback_info)) = player.playback_info().await else {
// If we're not playing anything, just wait until we are
return ControlFlow::Continue(());
};
if playback_info.track_id() != self.track {
// We're playing another track, reload the lyrics!
let lyrics = match player.get_lyrics().await {
Ok(lyrics) => lyrics,
Err(why) => {
error!("Failed to retrieve lyrics: {why}");
return ControlFlow::Break(());
}
};
self.lyrics = lyrics;
self.page = 0;
self.track = playback_info.track_id();
if let Err(why) = self
.message
.edit(
&self.ctx,
EditMessage::new()
.embed(lyrics_embed(&self.lyrics, &playback_info, self.page))
.components(vec![lyrics_buttons(
&self.guild_id,
&self.lyrics,
self.page,
)]),
)
.await
{
error!("Failed to update lyrics: {why}");
return ControlFlow::Break(());
}
return ControlFlow::Continue(());
}
// We're still playing the same song, check if we need to update the page
let Some(lyrics) = &self.lyrics else {
// No lyrics in current song, just continue until we have one with
return ControlFlow::Continue(());
};
if !matches!(lyrics.lyrics.sync_type, SyncType::LineSynced) {
// Only synced lyrics should auto-swap to new pages
return ControlFlow::Continue(());
}
let new_page = page_at_position(lyrics, playback_info.current_position()).unwrap_or(0);
if new_page != self.page {
// We've arrived on a new page: swap em up!
self.page = new_page;
if let Err(why) = self
.message
.edit(
&self.ctx,
EditMessage::new()
.embed(lyrics_embed(&self.lyrics, &playback_info, new_page))
.components(vec![lyrics_buttons(&self.guild_id, &self.lyrics, new_page)]),
)
.await
{
error!("Failed to update lyrics: {why}");
return ControlFlow::Break(());
}
}
ControlFlow::Continue(())
}
async fn handle_press(&mut self, press: ComponentInteraction) -> ControlFlow<(), ()> {
let next = match press.data.custom_id.split(':').nth(1) {
Some("next") => true,
Some("prev") => false,
_ => return ControlFlow::Continue(()),
};
let Some(lyrics) = &self.lyrics else {
return ControlFlow::Continue(());
};
if !matches!(lyrics.lyrics.sync_type, SyncType::Unsynced) {
// Only allow manual swapping if lyrics are unsynced
return ControlFlow::Continue(());
}
let length = lyrics
.lyrics
.lines
.iter()
.fold(0, |acc, line| acc + line.words.len());
let pages = length / PAGE_LENGTH + if length % PAGE_LENGTH > 0 { 1 } else { 0 };
let Ok(player) = self.session.player().await else {
return ControlFlow::Continue(());
};
let Ok(Some(playback_info)) = player.playback_info().await else {
return ControlFlow::Continue(());
};
match next {
true if self.page < pages - 1 => self.page += 1,
false if self.page > 0 => self.page -= 1,
_ => return ControlFlow::Continue(()),
}
if let Err(why) = self
.message
.edit(
&self.ctx,
EditMessage::new()
.embed(lyrics_embed(&self.lyrics, &playback_info, self.page))
.components(vec![lyrics_buttons(
&self.guild_id,
&self.lyrics,
self.page,
)]),
)
.await
{
error!("Failed to update lyrics: {why}");
return ControlFlow::Break(());
}
ControlFlow::Continue(())
}
}
async fn respond_not_playing(context: &Context, interaction: CommandInteraction) -> Result<()> {
interaction
.create_response(
context,
CreateInteractionResponse::Message(
CreateInteractionResponseMessage::new()
.embed(not_playing_embed())
.ephemeral(true),
),
)
.await?;
Ok(())
}
fn not_playing_embed() -> CreateEmbed {
CreateEmbed::new()
.title("Cannot get lyrics")
.description("I'm currently not playing any music in this server.")
.color(Colors::Error)
}
fn lyrics_embed(lyrics: &Option<Lyrics>, playback_info: &PlaybackInfo, page: usize) -> CreateEmbed {
match (lyrics, playback_info.artists()) {
(Some(lyrics), Some(artists)) => {
let length = lyrics
.lyrics
.lines
.iter()
.fold(0, |acc, line| acc + line.words.len());
let page = &into_pages(&lyrics.lyrics.lines)
[if page * PAGE_LENGTH > length { 0 } else { page }];
let title = format!(
"{} - {}",
playback_info.name(),
artists
.0
.into_iter()
.map(|artist| artist.name)
.collect::<Vec<_>>()
.join(", "),
);
let description = page
.iter()
.map(|page| page.words.replace('♪', "\n♪\n"))
.collect::<Vec<_>>()
.join("\n");
let mut footer = format!("Lyrics provided by {}", lyrics.lyrics.provider_display_name);
if matches!(lyrics.lyrics.sync_type, SyncType::LineSynced) {
footer.push_str(" | Synced to song");
}
CreateEmbed::new()
.title(title)
.description(description)
.footer(CreateEmbedFooter::new(footer))
.color(Colors::Info)
}
_ => CreateEmbed::new()
.title("No lyrics available")
.description("This current track has no lyrics available. Just enjoy the tunes!")
.color(Colors::Info),
}
}
fn lyrics_buttons(id: &str, lyrics: &Option<Lyrics>, page: usize) -> CreateActionRow {
let (can_prev, can_next) = match lyrics {
Some(lyrics) => match lyrics.lyrics.sync_type {
SyncType::Unsynced => {
// Only unsynced lyrics can have its pages flipped through by the user
let length = lyrics
.lyrics
.lines
.iter()
.fold(0, |acc, line| acc + line.words.len());
let pages = length / PAGE_LENGTH + if length % PAGE_LENGTH > 0 { 1 } else { 0 };
(page > 0, page < pages - 1)
}
SyncType::LineSynced => (false, false),
},
None => (false, false),
};
CreateActionRow::Buttons(vec![
CreateButton::new(format!("lyrics:prev:{id}"))
.disabled(!can_prev)
.label("<"),
CreateButton::new(format!("lyrics:next:{id}"))
.disabled(!can_next)
.label(">"),
])
}
fn into_pages(lines: &[Line]) -> Vec<Vec<Line>> {
let mut result = vec![];
let mut current = vec![];
let mut current_position = 0;
for line in lines {
if current_position + line.words.len() > PAGE_LENGTH {
result.push(current);
current = vec![line.clone()];
current_position = line.words.len();
continue;
}
current.push(line.clone());
current_position += line.words.len();
}
result.push(current);
result
}
fn page_at_position(lyrics: &Lyrics, position: u32) -> Option<usize> {
let pages = into_pages(&lyrics.lyrics.lines);
for (i, line) in pages.iter().enumerate() {
if let Some(first) = line.first() {
let Ok(time) = first
.start_time_ms
.parse::<u32>()
.map(|v| v.saturating_sub(TIME_OFFSET))
else {
return None;
};
if position < time {
return Some(if i == 0 { 0 } else { i - 1 });
}
}
if let (Some(first), Some(last)) = (line.first(), line.last()) {
let (Ok(first), Ok(last)) = (
first
.start_time_ms
.parse::<u32>()
.map(|v| v.saturating_sub(TIME_OFFSET)),
last.start_time_ms
.parse::<u32>()
.map(|v| v.saturating_sub(TIME_OFFSET)),
) else {
return None;
};
if position >= first && position <= last {
return Some(i);
}
}
}
Some(pages.len() - 1)
}
+126
View File
@@ -0,0 +1,126 @@
use anyhow::Result;
use serenity::all::{ChannelId, GuildId, UserId};
use songbird::Songbird;
use spoticord_database::Database;
use std::{
collections::HashMap,
sync::{Arc, Mutex},
};
use super::{Session, SessionHandle};
#[derive(Clone)]
pub struct SessionManager {
songbird: Arc<Songbird>,
database: Database,
sessions: Arc<Mutex<HashMap<GuildId, SessionHandle>>>,
owners: Arc<Mutex<HashMap<UserId, SessionHandle>>>,
}
pub enum SessionQuery {
Guild(GuildId),
Owner(UserId),
}
impl SessionManager {
pub fn new(songbird: Arc<Songbird>, database: Database) -> Self {
Self {
songbird,
database,
sessions: Arc::new(Mutex::new(HashMap::new())),
owners: Arc::new(Mutex::new(HashMap::new())),
}
}
pub async fn create_session(
&self,
context: &serenity::all::Context,
guild_id: GuildId,
voice_channel_id: ChannelId,
text_channel_id: ChannelId,
owner: UserId,
) -> Result<SessionHandle> {
let handle = Session::create(
self.clone(),
context,
guild_id,
voice_channel_id,
text_channel_id,
owner,
)
.await?;
self.sessions
.lock()
.expect("mutex poisoned")
.insert(guild_id, handle.clone());
self.owners
.lock()
.expect("mutex poisoned")
.insert(owner, handle.clone());
Ok(handle)
}
pub fn get_session(&self, query: SessionQuery) -> Option<SessionHandle> {
match query {
SessionQuery::Guild(guild) => self
.sessions
.lock()
.expect("mutex poisoned")
.get(&guild)
.cloned(),
SessionQuery::Owner(owner) => self
.owners
.lock()
.expect("mutex poisoned")
.get(&owner)
.cloned(),
}
}
pub fn remove_session(&self, query: SessionQuery) {
match query {
SessionQuery::Guild(guild) => {
self.sessions.lock().expect("mutex poisoned").remove(&guild)
}
SessionQuery::Owner(owner) => {
self.owners.lock().expect("mutex poisoned").remove(&owner)
}
};
}
pub fn get_all_sessions(&self) -> Vec<SessionHandle> {
self.sessions
.lock()
.expect("mutex poisoned")
.values()
.cloned()
.collect()
}
/// Disconnects all active sessions and clears out all handles.
///
/// The session manager can still create new sessions after all sessions have been shut down.
/// Sessions might still be created during shutdown.
pub async fn shutdown_all(&self) {
let sessions = self.get_all_sessions();
for session in sessions {
session.disconnect().await;
}
self.owners.lock().expect("mutex poisoned").clear();
self.sessions.lock().expect("mutex poisoned").clear();
}
pub fn songbird(&self) -> Arc<Songbird> {
self.songbird.clone()
}
pub fn database(&self) -> Database {
self.database.clone()
}
}
+401
View File
@@ -0,0 +1,401 @@
use anyhow::{anyhow, Result};
use log::{error, trace};
use serenity::{
all::{
ButtonStyle, CommandInteraction, ComponentInteraction, ComponentInteractionCollector,
Context, CreateActionRow, CreateButton, CreateEmbed, CreateEmbedAuthor, CreateEmbedFooter,
CreateInteractionResponse, CreateInteractionResponseFollowup,
CreateInteractionResponseMessage, EditMessage, Message, User,
},
futures::StreamExt,
};
use spoticord_player::{info::PlaybackInfo, PlayerHandle};
use spoticord_utils::discord::Colors;
use std::{ops::ControlFlow, time::Duration};
use tokio::{sync::mpsc, time::Instant};
use crate::{Session, SessionHandle};
#[derive(Debug)]
pub enum Command {
InvokeUpdate,
}
pub struct PlaybackEmbed {
id: u64,
ctx: Context,
session: SessionHandle,
message: Message,
last_update: Instant,
update_in: Option<Duration>,
rx: mpsc::Receiver<Command>,
}
impl PlaybackEmbed {
pub async fn create(
session: &Session,
handle: SessionHandle,
interaction: CommandInteraction,
) -> Result<Option<PlaybackEmbedHandle>> {
let ctx = session.context.clone();
if !session.active {
respond_not_playing(&ctx, interaction).await?;
return Ok(None);
}
let owner = session.owner.to_user(&ctx).await?;
let Some(playback_info) = session.player.playback_info().await? else {
respond_not_playing(&ctx, interaction).await?;
return Ok(None);
};
let ctx_id = interaction.id.get();
// Send initial reply
interaction
.create_response(
&ctx,
CreateInteractionResponse::Message(
CreateInteractionResponseMessage::new()
.embed(build_embed(&playback_info, &owner))
.components(vec![build_buttons(ctx_id, playback_info.playing())]),
),
)
.await?;
// Retrieve message instead of editing interaction response, as those tokens are only valid for 15 minutes
let message = interaction.get_response(&ctx).await?;
let collector = ComponentInteractionCollector::new(&ctx)
.filter(move |press| press.data.custom_id.starts_with(&ctx_id.to_string()))
.timeout(Duration::from_secs(3600 * 24));
let (tx, rx) = mpsc::channel(16);
let this = Self {
id: ctx_id,
ctx,
session: handle,
message,
last_update: Instant::now(),
update_in: None,
rx,
};
tokio::spawn(this.run(collector));
Ok(Some(PlaybackEmbedHandle { tx }))
}
async fn run(mut self, collector: ComponentInteractionCollector) {
let mut stream = collector.stream();
loop {
tokio::select! {
opt_command = self.rx.recv() => {
let Some(command) = opt_command else {
break;
};
if self.handle_command(command).await.is_break() {
break;
}
},
opt_press = stream.next() => {
let Some(press) = opt_press else {
break;
};
self.handle_press(press).await;
}
_ = async {
if let Some(update_in) = self.update_in.take()
{
tokio::time::sleep(update_in).await;
}
}, if self.update_in.is_some() => {
if self.update_embed().await.is_break() {
break;
}
}
}
}
}
async fn handle_command(&mut self, command: Command) -> ControlFlow<(), ()> {
trace!("Received command: {command:?}");
match command {
Command::InvokeUpdate => {
if self.last_update.elapsed() < Duration::from_secs(2) {
if self.update_in.is_some() {
return ControlFlow::Continue(());
}
self.update_in = Some(Duration::from_secs(2) - self.last_update.elapsed());
} else {
self.update_embed().await?;
}
}
}
ControlFlow::Continue(())
}
async fn handle_press(&self, press: ComponentInteraction) {
trace!("Received button press: {press:?}");
let Ok((player, playback_info, owner)) = self.get_info().await else {
_ = press
.create_followup(
&self.ctx,
CreateInteractionResponseFollowup::new()
.embed(
CreateEmbed::new()
.title("Cannot perform action")
.description("I'm currently not playing any music in this server"),
)
.ephemeral(true),
)
.await;
return;
};
if press.user.id != owner.id {
_ = press
.create_followup(
&self.ctx,
CreateInteractionResponseFollowup::new()
.embed(
CreateEmbed::new()
.title("Cannot perform action")
.description("Only the host may use the media buttons"),
)
.ephemeral(true),
)
.await;
return;
}
match press.data.custom_id.split('-').last() {
Some("next") => player.next_track().await,
Some("prev") => player.previous_track().await,
Some("pause") => {
if playback_info.playing() {
player.pause().await
} else {
player.play().await
}
}
_ => {}
}
_ = press
.create_response(&self.ctx, CreateInteractionResponse::Acknowledge)
.await;
}
async fn get_info(&self) -> Result<(PlayerHandle, PlaybackInfo, User)> {
let player = self.session.player().await?;
let owner = self.session.owner().await?.to_user(&self.ctx).await?;
let playback_info = player
.playback_info()
.await?
.ok_or_else(|| anyhow!("No playback info present"))?;
Ok((player, playback_info, owner))
}
async fn update_embed(&mut self) -> ControlFlow<(), ()> {
self.update_in = None;
let Ok(owner) = self.session.owner().await else {
_ = self.update_not_playing().await;
return ControlFlow::Break(());
};
let Ok(player) = self.session.player().await else {
_ = self.update_not_playing().await;
return ControlFlow::Break(());
};
let Ok(Some(playback_info)) = player.playback_info().await else {
_ = self.update_not_playing().await;
return ControlFlow::Break(());
};
let owner = match owner.to_user(&self.ctx).await {
Ok(owner) => owner,
Err(why) => {
error!("Failed to resolve owner: {why}");
return ControlFlow::Break(());
}
};
if let Err(why) = self
.message
.edit(
&self.ctx,
EditMessage::new()
.embed(build_embed(&playback_info, &owner))
.components(vec![build_buttons(self.id, playback_info.playing())]),
)
.await
{
error!("Failed to update playback embed: {why}");
return ControlFlow::Break(());
};
self.last_update = Instant::now();
ControlFlow::Continue(())
}
async fn update_not_playing(&mut self) -> Result<()> {
self.message
.edit(&self.ctx, EditMessage::new().embed(not_playing_embed()))
.await?;
Ok(())
}
}
pub struct PlaybackEmbedHandle {
tx: mpsc::Sender<Command>,
}
impl PlaybackEmbedHandle {
pub fn is_valid(&self) -> bool {
!self.tx.is_closed()
}
pub async fn invoke_update(&self) -> Result<()> {
self.tx.send(Command::InvokeUpdate).await?;
Ok(())
}
}
async fn respond_not_playing(context: &Context, interaction: CommandInteraction) -> Result<()> {
interaction
.create_response(
context,
CreateInteractionResponse::Message(
CreateInteractionResponseMessage::new()
.embed(not_playing_embed())
.ephemeral(true),
),
)
.await?;
Ok(())
}
fn not_playing_embed() -> CreateEmbed {
CreateEmbed::new()
.title("Cannot display song details")
.description("I'm currently not playing any music in this server.")
.color(Colors::Error)
}
fn build_embed(playback_info: &PlaybackInfo, owner: &User) -> CreateEmbed {
let mut description = String::new();
description += &format!("## [{}]({})\n", playback_info.name(), playback_info.url());
if let Some(artists) = playback_info.artists() {
let artists = artists
.iter()
.map(|artist| {
format!(
"[{}](https://open.spotify.com/artist/{})",
artist.name,
artist.id.to_base62().expect("invalid artist")
)
})
.collect::<Vec<_>>()
.join(", ");
description += &format!("By {artists}\n\n");
}
if let Some(show_name) = playback_info.show_name() {
description += &format!("On {show_name}\n\n");
}
let position = playback_info.current_position();
let index = position * 20 / playback_info.duration();
description.push_str(if playback_info.playing() {
"▶️ "
} else {
"⏸️ "
});
for i in 0..20 {
if i == index {
description.push('🔵');
} else {
description.push('▬');
}
}
description.push_str("\n:alarm_clock: ");
description.push_str(&format!(
"{} / {}",
spoticord_utils::time_to_string(position / 1000),
spoticord_utils::time_to_string(playback_info.duration() / 1000)
));
CreateEmbed::new()
.author(
CreateEmbedAuthor::new("Currently Playing")
.icon_url("https://spoticord.com/spotify-logo.png"),
)
.description(description)
.thumbnail(playback_info.thumbnail())
.footer(
CreateEmbedFooter::new(owner.global_name.as_ref().unwrap_or(&owner.name))
.icon_url(owner.face()),
)
.color(Colors::Info)
}
fn build_buttons(id: u64, playing: bool) -> CreateActionRow {
let prev_button_id = format!("{id}-prev");
let next_button_id = format!("{id}-next");
let pause_button_id = format!("{id}-pause");
let prev_button = CreateButton::new(prev_button_id)
.style(ButtonStyle::Primary)
.label("<<");
let next_button = CreateButton::new(next_button_id)
.style(ButtonStyle::Primary)
.label(">>");
let pause_button = CreateButton::new(pause_button_id)
.style(if playing {
ButtonStyle::Danger
} else {
ButtonStyle::Success
})
.label(if playing { "Pause" } else { "Play" });
CreateActionRow::Buttons(vec![prev_button, pause_button, next_button])
}