diff --git a/Cargo.lock b/Cargo.lock index 0f81a07..7b203e0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -454,6 +454,7 @@ version = "0.1.0" dependencies = [ "crabidy-core", "crossterm", + "dirs", "flume", "notify-rust", "ratatui", @@ -461,6 +462,9 @@ dependencies = [ "tokio", "tokio-stream", "tonic", + "tracing", + "tracing-appender", + "tracing-subscriber", ] [[package]] diff --git a/cbd-tui/Cargo.toml b/cbd-tui/Cargo.toml index c8afea8..59d595a 100644 --- a/cbd-tui/Cargo.toml +++ b/cbd-tui/Cargo.toml @@ -6,6 +6,7 @@ edition.workspace = true [dependencies] crabidy-core.workspace = true crossterm.workspace = true +dirs.workspace = true flume.workspace = true notify-rust.workspace = true ratatui.workspace = true @@ -13,3 +14,6 @@ serde.workspace = true tokio = { workspace = true, features = ["full"] } tokio-stream.workspace = true tonic.workspace = true +tracing.workspace = true +tracing-appender.workspace = true +tracing-subscriber.workspace = true diff --git a/cbd-tui/src/app/now_playing.rs b/cbd-tui/src/app/now_playing.rs index e04a87b..e94d325 100644 --- a/cbd-tui/src/app/now_playing.rs +++ b/cbd-tui/src/app/now_playing.rs @@ -56,11 +56,14 @@ impl NowPlaying { } else { format!("{} by {}", track.title, track.artist,) }; - Notification::new() + // A missing notification daemon must not crash the TUI. + if let Err(err) = Notification::new() .summary("Now playing") .body(&body) .show() - .unwrap(); + { + tracing::debug!("could not show desktop notification: {err}"); + } } self.track = active; } diff --git a/cbd-tui/src/main.rs b/cbd-tui/src/main.rs index c906644..3d1c034 100644 --- a/cbd-tui/src/main.rs +++ b/cbd-tui/src/main.rs @@ -27,11 +27,45 @@ use tokio_stream::StreamExt; use app::{App, MessageFromUi, MessageToUi, StatefulList, UiFocus}; use config::Config; use rpc::RpcClient; +use tracing::{error, info, warn}; static CONFIG: OnceLock = OnceLock::new(); +/// Logs to a file: the terminal is owned by the TUI, so writing log lines to +/// stdout/stderr would corrupt the interface. +fn init_tracing() -> Option { + use tracing_subscriber::{prelude::*, EnvFilter}; + + let log_dir = dirs::state_dir() + .or_else(dirs::cache_dir) + .unwrap_or_else(std::env::temp_dir) + .join("crabidy"); + if let Err(err) = std::fs::create_dir_all(&log_dir) { + eprintln!( + "could not create log directory {}: {err}", + log_dir.display() + ); + return None; + } + let file_appender = tracing_appender::rolling::daily(&log_dir, "cbd-tui.log"); + let (non_blocking, guard) = tracing_appender::non_blocking(file_appender); + let env_filter = EnvFilter::try_from_default_env() + .unwrap_or_else(|_| EnvFilter::new("info,cbd_tui=debug,crabidy_core=debug")); + tracing_subscriber::registry() + .with(env_filter) + .with( + tracing_subscriber::fmt::layer() + .with_writer(non_blocking) + .with_ansi(false) + .with_target(true), + ) + .init(); + Some(guard) +} + #[tokio::main] async fn main() -> Result<(), Box> { + let _log_guard = init_tracing(); let config = CONFIG.get_or_init(|| crabidy_core::init_config("cbd-tui.toml")); let (ui_tx, rx): (Sender, Receiver) = flume::unbounded(); @@ -52,6 +86,7 @@ async fn orchestrate<'a>( config: &'static Config, (tx, rx): (Sender, Receiver), ) -> Result<(), Box> { + info!(address = config.server.address, "connecting to server"); let mut rpc_client = rpc::RpcClient::connect(&config.server.address).await?; if let Some(root_node) = rpc_client.get_library_node("node:/").await? { @@ -59,11 +94,12 @@ async fn orchestrate<'a>( } let init_data = rpc_client.init().await?; + info!("received initial state from server"); tx.send_async(MessageToUi::Init(init_data)).await?; loop { - if let Err(er) = poll(&mut rpc_client, &rx, &tx).await { - println!("ERROR"); + if let Err(err) = poll(&mut rpc_client, &rx, &tx).await { + error!("request to server failed: {err}"); } } } @@ -135,8 +171,10 @@ async fn poll( tx.send_async(MessageToUi::Update(update)).await?; } } - Err(_) => { + Err(err) => { + warn!("update stream broke, reconnecting: {err}"); rpc_client.reconnect_update_stream().await; + info!("update stream reconnected"); } } diff --git a/crabidy-core/src/lib.rs b/crabidy-core/src/lib.rs index f9ce8ca..94c0ece 100644 --- a/crabidy-core/src/lib.rs +++ b/crabidy-core/src/lib.rs @@ -39,6 +39,8 @@ impl std::fmt::Display for ProviderError { } } +impl std::error::Error for ProviderError {} + impl LibraryNode { pub fn new() -> Self { Self { diff --git a/crabidy-server/src/lib.rs b/crabidy-server/src/lib.rs index 0396324..88d7415 100644 --- a/crabidy-server/src/lib.rs +++ b/crabidy-server/src/lib.rs @@ -16,10 +16,11 @@ pub struct QueueManager { impl From for Queue { fn from(queue_manager: QueueManager) -> Self { Self { + // A clock step backwards must not panic the playback loop. timestamp: queue_manager .created_at .elapsed() - .expect("failed to get elapsed time") + .unwrap_or_default() .as_secs(), current_position: queue_manager.current_position() as u32, tracks: queue_manager.tracks, diff --git a/crabidy-server/src/main.rs b/crabidy-server/src/main.rs index 8997d38..e1ba0cc 100644 --- a/crabidy-server/src/main.rs +++ b/crabidy-server/src/main.rs @@ -3,8 +3,8 @@ use crabidy_core::proto::crabidy::{ crabidy_service_server::CrabidyServiceServer, InitResponse, LibraryNode, PlayState, Track, }; use crabidy_core::{ProviderClient, ProviderError}; -use tracing::{debug_span, error, info, instrument, level_filters, warn, Span}; -use tracing_subscriber::{filter::Targets, prelude::*}; +use tracing::{debug, error, info, instrument, warn, Span}; +use tracing_subscriber::{prelude::*, EnvFilter}; mod playback; use playback::Playback; @@ -15,34 +15,17 @@ use rpc::RpcService; use tonic::{transport::Server, Result}; +const LISTEN_ADDR: &str = "0.0.0.0:50051"; + #[tokio::main] async fn main() -> Result<(), Box> { - if let Err(err) = tracing_log::LogTracer::init_with_filter(log::LevelFilter::Debug) { - println!("Failed to initialize log tracer: {}", err); - } - let (non_blocking, _guard) = tracing_appender::non_blocking(std::io::stderr()); - - let targets_filter = Targets::new() - .with_target("crabidy_server", tracing::level_filters::LevelFilter::DEBUG) - .with_target("tidaldy", level_filters::LevelFilter::DEBUG); - let subscriber = tracing_subscriber::fmt::layer() - .with_writer(non_blocking) - .with_file(true) - .with_line_number(true); - - let registry = tracing_subscriber::registry() - .with(targets_filter) - .with(subscriber); - - tracing::subscriber::set_global_default(registry) - .expect("Setting the default tracing subscriber failed"); - - info!("audio player started initialized"); + let _log_guard = init_tracing(); let (update_tx, _) = tokio::sync::broadcast::channel(2048); - let orchestrator = ProviderOrchestrator::init("") - .await - .expect("failed to init orchestrator"); + let orchestrator = ProviderOrchestrator::init("").await.map_err(|err| { + error!("failed to init provider orchestrator: {err}"); + err + })?; let playback = Playback::new(update_tx.clone(), orchestrator.provider_tx.clone()); @@ -52,7 +35,7 @@ async fn main() -> Result<(), Box> { std::thread::spawn(|| { poll_play_bus(player_msg, playback_tx); }); - info!("gstreamer bus handler started"); + info!("player message forwarder started"); let crabidy_service = RpcService::new( update_tx, @@ -64,7 +47,8 @@ async fn main() -> Result<(), Box> { playback.run(); info!("playback started"); - let addr = "0.0.0.0:50051".parse()?; + let addr = LISTEN_ADDR.parse()?; + info!(%addr, "grpc server listening"); Server::builder() .add_service(CrabidyServiceServer::new(crabidy_service)) .serve(addr) @@ -73,165 +57,216 @@ async fn main() -> Result<(), Box> { Ok(()) } +/// Installs the global tracing subscriber. +/// +/// The filter honors `RUST_LOG`; without it, our own crates log at debug and +/// everything else at info. Returns the guard that flushes the non-blocking +/// writer on shutdown. +fn init_tracing() -> tracing_appender::non_blocking::WorkerGuard { + let (non_blocking, guard) = tracing_appender::non_blocking(std::io::stderr()); + + let env_filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| { + EnvFilter::new( + "info,crabidy_server=debug,crabidy_core=debug,tidaldy=debug,audio_player=debug", + ) + }); + + let fmt_layer = tracing_subscriber::fmt::layer() + .with_writer(non_blocking) + .with_target(true) + .with_file(true) + .with_line_number(true); + + tracing_subscriber::registry() + .with(env_filter) + .with(fmt_layer) + .init(); + + // Forward records from libraries that use the `log` crate (symphonia, + // cpal, ...) into tracing. + if let Err(err) = tracing_log::LogTracer::init() { + warn!("failed to initialize log-to-tracing bridge: {err}"); + } + guard +} + +/// Forwards player engine events into the playback message loop. #[instrument(skip(rx, tx))] fn poll_play_bus(rx: flume::Receiver, tx: flume::Sender) { for msg in rx.iter() { - let span = debug_span!("play-chan"); - match msg { + let command = match msg { PlayerMessage::EndOfStream => { - if let Err(err) = tx.send(PlaybackMessage::Next { span }) { - error!("failed to send next message: {}", err); - } - } - PlayerMessage::Stopped => { - if let Err(err) = tx.send(PlaybackMessage::StateChanged { - state: PlayState::Stopped, - span, - }) { - error!("failed to send stopped message: {}", err); - } - } - PlayerMessage::Paused => { - if let Err(err) = tx.send(PlaybackMessage::StateChanged { - state: PlayState::Paused, - span, - }) { - error!("failed to send paused message: {}", err); - } - } - PlayerMessage::Playing => { - if let Err(err) = tx.send(PlaybackMessage::StateChanged { - state: PlayState::Playing, - span, - }) { - error!("failed to send playing message: {}", err); - } - } - PlayerMessage::Elapsed { duration, elapsed } => { - if let Err(err) = tx.send(PlaybackMessage::PostitionChanged { - duration: duration.as_millis() as u32, - position: elapsed.as_millis() as u32, - span, - }) { - error!("failed to send elapsed message: {}", err); - } - } - PlayerMessage::Duration { duration } => { - if let Err(err) = tx.send(PlaybackMessage::PostitionChanged { - duration: duration.as_millis() as u32, - position: 0, - span, - }) { - error!("failed to send duration message: {}", err); - } + debug!("player reported end of stream"); + PlaybackCommand::Next } + PlayerMessage::Stopped => PlaybackCommand::StateChanged { + state: PlayState::Stopped, + }, + PlayerMessage::Paused => PlaybackCommand::StateChanged { + state: PlayState::Paused, + }, + PlayerMessage::Playing => PlaybackCommand::StateChanged { + state: PlayState::Playing, + }, + PlayerMessage::Elapsed { duration, elapsed } => PlaybackCommand::PositionChanged { + duration: duration.as_millis() as u32, + position: elapsed.as_millis() as u32, + }, + PlayerMessage::Duration { duration } => PlaybackCommand::PositionChanged { + duration: duration.as_millis() as u32, + position: 0, + }, + }; + if let Err(err) = tx.send(PlaybackMessage::new(command)) { + error!("failed to forward player message: {err}"); + return; + } + } + warn!("player message channel closed"); +} + +/// A command for the provider orchestrator, tagged with the tracing span that +/// was current when it was sent so the handler can attribute its events to +/// the originating request. +#[derive(Debug)] +pub struct ProviderMessage { + pub span: Span, + pub command: ProviderCommand, +} + +impl ProviderMessage { + pub fn new(command: ProviderCommand) -> Self { + Self { + span: Span::current(), + command, } } } #[derive(Debug)] -pub enum ProviderMessage { +pub enum ProviderCommand { GetLibraryNode { uuid: String, result_tx: flume::Sender>, - span: Span, }, GetTrack { uuid: String, result_tx: flume::Sender>, - span: Span, }, GetTrackUrls { uuid: String, result_tx: flume::Sender, ProviderError>>, - span: Span, }, FlattenNode { uuid: String, result_tx: flume::Sender>, - span: Span, }, } +impl ProviderCommand { + pub fn name(&self) -> &'static str { + match self { + Self::GetLibraryNode { .. } => "get_library_node", + Self::GetTrack { .. } => "get_track", + Self::GetTrackUrls { .. } => "get_track_urls", + Self::FlattenNode { .. } => "flatten_node", + } + } +} + +/// A command for the playback loop, tagged like [`ProviderMessage`]. #[derive(Debug)] -pub enum PlaybackMessage { +pub struct PlaybackMessage { + pub span: Span, + pub command: PlaybackCommand, +} + +impl PlaybackMessage { + pub fn new(command: PlaybackCommand) -> Self { + Self { + span: Span::current(), + command, + } + } +} + +#[derive(Debug)] +pub enum PlaybackCommand { Init { result_tx: flume::Sender, - span: Span, }, Replace { uuids: Vec, - span: Span, }, Queue { uuids: Vec, - span: Span, }, Append { uuids: Vec, - span: Span, }, Remove { positions: Vec, - span: Span, }, Insert { position: u32, uuids: Vec, - span: Span, }, Clear { exclude_current: bool, - span: Span, }, SetCurrent { position: u32, - span: Span, - }, - ToggleShuffle { - span: Span, - }, - ToggleRepeat { - span: Span, - }, - TogglePlay { - span: Span, - }, - Stop { - span: Span, }, + ToggleShuffle, + ToggleRepeat, + TogglePlay, + Stop, ChangeVolume { delta: f32, - span: Span, - }, - ToggleMute { - span: Span, - }, - Next { - span: Span, - }, - Prev { - span: Span, - }, - RestartTrack { - span: Span, }, + ToggleMute, + Next, + Prev, + RestartTrack, StateChanged { state: PlayState, - span: Span, }, VolumeChanged { volume: f32, - span: Span, }, - MuteChanged { muted: bool, - span: Span, }, - PostitionChanged { + PositionChanged { duration: u32, position: u32, - span: Span, }, } + +impl PlaybackCommand { + pub fn name(&self) -> &'static str { + match self { + Self::Init { .. } => "init", + Self::Replace { .. } => "replace", + Self::Queue { .. } => "queue", + Self::Append { .. } => "append", + Self::Remove { .. } => "remove", + Self::Insert { .. } => "insert", + Self::Clear { .. } => "clear", + Self::SetCurrent { .. } => "set_current", + Self::ToggleShuffle => "toggle_shuffle", + Self::ToggleRepeat => "toggle_repeat", + Self::TogglePlay => "toggle_play", + Self::Stop => "stop", + Self::ChangeVolume { .. } => "change_volume", + Self::ToggleMute => "toggle_mute", + Self::Next => "next", + Self::Prev => "prev", + Self::RestartTrack => "restart_track", + Self::StateChanged { .. } => "state_changed", + Self::VolumeChanged { .. } => "volume_changed", + Self::MuteChanged { .. } => "mute_changed", + Self::PositionChanged { .. } => "position_changed", + } + } +} diff --git a/crabidy-server/src/playback.rs b/crabidy-server/src/playback.rs index a2fdd4a..aab26d8 100644 --- a/crabidy-server/src/playback.rs +++ b/crabidy-server/src/playback.rs @@ -1,5 +1,4 @@ -use crate::PlaybackMessage; -use crate::ProviderMessage; +use crate::{PlaybackCommand, PlaybackMessage, ProviderCommand, ProviderMessage}; use audio_player::Player; use crabidy_core::proto::crabidy::QueueModifiers; use crabidy_core::proto::crabidy::{ @@ -9,8 +8,7 @@ use crabidy_core::proto::crabidy::{ use crabidy_core::ProviderError; use crabidy_server::QueueManager; use std::sync::Mutex; -use tracing::debug_span; -use tracing::{debug, error, instrument, trace, warn, Instrument}; +use tracing::{debug, debug_span, error, instrument, trace, warn, Instrument}; pub struct Playback { update_tx: tokio::sync::broadcast::Sender, @@ -27,7 +25,7 @@ impl Playback { update_tx: tokio::sync::broadcast::Sender, provider_tx: flume::Sender, ) -> Self { - let (playback_tx, playback_rx) = flume::bounded(10); + let (playback_tx, playback_rx) = flume::bounded(64); let queue = Mutex::new(QueueManager::new()); let state = Mutex::new(PlayState::Stopped); let player = Player::default(); @@ -44,633 +42,465 @@ impl Playback { pub fn run(self) { tokio::spawn(async move { - while let Ok(message) = self.playback_rx.recv_async().await { - match message { - PlaybackMessage::Init { result_tx, span } => { - let _e = span.enter(); - let repeat; - let shuffle; - let response = { - let Ok(queue) = self.queue.lock() else { - error!("failed to get queue lock"); - continue; - }; - debug!("got queue lock"); - repeat = queue.repeat; - shuffle = queue.shuffle; - let queue_track = QueueTrack { - queue_position: queue.current_position() as u32, - track: queue.current_track(), - }; - trace!("queue_track {:?}", queue_track); - debug!("released queue_track lock"); + while let Ok(PlaybackMessage { span, command }) = self.playback_rx.recv_async().await { + // Attribute all handler events to a span that is a child of + // the span that was current when the command was sent. + let handler_span = + debug_span!(parent: &span, "playback_command", command = command.name()); + self.handle_command(command).instrument(handler_span).await; + } + warn!("playback message channel closed, loop exiting"); + }); + } - let position = TrackPosition { - duration: 0, - position: 0, - }; - trace!("position {:?}", position); - let play_state = { - debug!("getting play state lock"); - let Ok(play_state) = self.state.lock() else { - error!("failed to get play state lock"); - continue; - }; - *play_state - }; - trace!("play_state {:?}", play_state); - debug!("released play state lock"); - InitResponse { - queue: Some(queue.clone().into()), - queue_track: Some(queue_track), - play_state: play_state as i32, - volume: 0.0, - mute: false, - position: Some(position), - mods: Some(QueueModifiers { repeat, shuffle }), - } + async fn handle_command(&self, command: PlaybackCommand) { + trace!("handling playback command"); + match command { + PlaybackCommand::Init { result_tx } => { + let response = { + let Ok(queue) = self.queue.lock() else { + error!("queue lock poisoned"); + return; + }; + let queue_track = QueueTrack { + queue_position: queue.current_position() as u32, + track: queue.current_track(), + }; + let position = TrackPosition { + duration: 0, + position: 0, + }; + let play_state = { + let Ok(play_state) = self.state.lock() else { + error!("play state lock poisoned"); + return; }; - trace!("response {:?}", response); - if let Err(err) = result_tx.send(response) { - error!("failed to send response: {:#?}", err); - } - } - PlaybackMessage::Replace { uuids, span } => { - let _e = span.enter(); - let mut all_tracks = Vec::new(); - for uuid in uuids { - if is_track(&uuid) { - if let Ok(track) = self.get_track(&uuid).in_current_span().await { - all_tracks.push(track); - } - } else { - let tracks = self.flatten_node(&uuid).in_current_span().await; - all_tracks.extend(tracks); - } - debug!("uuid: {:?}", uuid); - } - trace!("got tracks {:?}", all_tracks); - let current = { - let Ok(mut queue) = self.queue.lock() else { - debug!("got queue lock"); - continue; - }; - debug!("got queue lock"); - queue.replace_with_tracks(&all_tracks); - let queue_update_tx = self.update_tx.clone(); - let update = StreamUpdate::Queue(queue.clone().into()); - if let Err(err) = queue_update_tx.send(update) { - trace!("{:?}", err) - }; - queue.current_track() - }; - debug!("got current {:?}", current); - self.play(current).in_current_span().await; + *play_state + }; + InitResponse { + queue: Some(queue.clone().into()), + queue_track: Some(queue_track), + play_state: play_state as i32, + volume: 0.0, + mute: false, + position: Some(position), + mods: Some(QueueModifiers { + repeat: queue.repeat, + shuffle: queue.shuffle, + }), } + }; + trace!(?response, "sending init response"); + if let Err(err) = result_tx.send(response) { + error!("failed to send init response: {err}"); + } + } - PlaybackMessage::Queue { uuids, span } => { - let _e = span.enter(); - debug!("queing"); - let mut all_tracks = Vec::new(); - for uuid in uuids { - if is_track(&uuid) { - if let Ok(track) = self.get_track(&uuid).in_current_span().await { - all_tracks.push(track); - } - } else { - let tracks = self.flatten_node(&uuid).in_current_span().await; - all_tracks.extend(tracks); - } - } - trace!("got tracks {:?}", all_tracks); - let track = { - let Ok(mut queue) = self.queue.lock() else { - debug!("got queue lock"); - continue; - }; - debug!("got queue lock"); - let track = queue.queue_tracks(&all_tracks); - let queue_update_tx = self.update_tx.clone(); - let update = StreamUpdate::Queue(queue.clone().into()); - if let Err(err) = queue_update_tx.send(update) { - trace!("{:?}", err) - } - track - }; - debug!("que lock released"); - self.play(track).in_current_span().await; - } + PlaybackCommand::Replace { uuids } => { + let all_tracks = self.resolve_tracks(uuids).await; + debug!(count = all_tracks.len(), "replacing queue"); + let current = { + let Ok(mut queue) = self.queue.lock() else { + error!("queue lock poisoned"); + return; + }; + queue.replace_with_tracks(&all_tracks); + self.broadcast(StreamUpdate::Queue(queue.clone().into())); + queue.current_track() + }; + self.play(current).await; + } - PlaybackMessage::Append { uuids, span } => { - let _e = span.enter(); - debug!("appending"); - let mut all_tracks = Vec::new(); - for uuid in uuids { - if is_track(&uuid) { - if let Ok(track) = self.get_track(&uuid).in_current_span().await { - all_tracks.push(track); - } - } else { - let tracks = self.flatten_node(&uuid).in_current_span().await; - all_tracks.extend(tracks); - } - } - trace!("got tracks {:?}", all_tracks); - let track = { - let Ok(mut queue) = self.queue.lock() else { - debug!("got queue lock"); - continue; - }; - debug!("got queue lock"); - let track = queue.append_tracks(&all_tracks); - let queue_update_tx = self.update_tx.clone(); - let update = StreamUpdate::Queue(queue.clone().into()); - if let Err(err) = queue_update_tx.send(update) { - trace!("{:?}", err) - } - track - }; - debug!("queue lock released"); - self.play(track).in_current_span().await; - } + PlaybackCommand::Queue { uuids } => { + let all_tracks = self.resolve_tracks(uuids).await; + debug!(count = all_tracks.len(), "queueing after current"); + let track = { + let Ok(mut queue) = self.queue.lock() else { + error!("queue lock poisoned"); + return; + }; + let track = queue.queue_tracks(&all_tracks); + self.broadcast(StreamUpdate::Queue(queue.clone().into())); + track + }; + self.play_if_some(track).await; + } - PlaybackMessage::Remove { positions, span } => { - let _e = span.enter(); - let is_last; - debug!("removing"); - let track = { - let Ok(mut queue) = self.queue.lock() else { - debug!("got queue lock"); - continue; - }; - debug!("got queue lock"); - is_last = queue.is_last_track(); - let track = queue.remove_tracks(&positions); - let queue_update_tx = self.update_tx.clone(); - let update = StreamUpdate::Queue(queue.clone().into()); - if let Err(err) = queue_update_tx.send(update) { - trace!("{:?}", err) - }; - track - }; - debug!("queue lock released"); - let state = { - let Ok(state) = self.state.lock() else { - error!("failed to get play state lock"); - continue; - }; - *state - }; - if state == PlayState::Playing { - if is_last { - if let Err(err) = self.player.stop().in_current_span().await { - error!("{:?}", err) - } - } else { - self.play(track).in_current_span().await; - } - } - } + PlaybackCommand::Append { uuids } => { + let all_tracks = self.resolve_tracks(uuids).await; + debug!(count = all_tracks.len(), "appending to queue"); + let track = { + let Ok(mut queue) = self.queue.lock() else { + error!("queue lock poisoned"); + return; + }; + let track = queue.append_tracks(&all_tracks); + self.broadcast(StreamUpdate::Queue(queue.clone().into())); + track + }; + self.play_if_some(track).await; + } - PlaybackMessage::Insert { - position, - uuids, - span, - } => { - let _e = span.enter(); - debug!("inserting"); - let mut all_tracks = Vec::new(); - for uuid in uuids { - if is_track(&uuid) { - if let Ok(track) = self.get_track(&uuid).in_current_span().await { - all_tracks.push(track); - } - } else { - let tracks = self.flatten_node(&uuid).in_current_span().await; - all_tracks.extend(tracks); - } - } - trace!("got tracks {:?}", all_tracks); - let track = { - let Ok(mut queue) = self.queue.lock() else { - debug!("got queue lock"); - continue; - }; - debug!("got queue lock"); - let track = queue.insert_tracks(position, &all_tracks); - let queue_update_tx = self.update_tx.clone(); - let update = StreamUpdate::Queue(queue.clone().into()); - if let Err(err) = queue_update_tx.send(update) { - trace!("{:?}", err) - }; - track - }; - debug!("queue lock released"); - self.play(track).in_current_span().await; - } - - PlaybackMessage::Clear { - exclude_current, - span, - } => { - let _e = span.enter(); - debug!("clearing"); - let should_stop = { - let Ok(mut queue) = self.queue.lock() else { - debug!("got queue lock"); - continue; - }; - debug!("got queue lock"); - let should_stop = queue.clear(exclude_current); - let queue_update_tx = self.update_tx.clone(); - let update = StreamUpdate::Queue(queue.clone().into()); - if let Err(err) = queue_update_tx.send(update) { - trace!("{:?}", err) - }; - should_stop - }; - debug!("queue lock released"); - if should_stop { - if let Err(err) = self.player.stop().in_current_span().await { - error!("{:?}", err) - } - } - } - - PlaybackMessage::SetCurrent { - position: queue_position, - span, - } => { - let _e = span.enter(); - debug!("setting current"); - let track = { - let Ok(mut queue) = self.queue.lock() else { - debug!("got queue lock"); - continue; - }; - debug!("got queue lock"); - queue.set_current_position(queue_position); - queue.current_track() - }; - debug!("quue lock released and got current {:?}", track); - self.play(track).in_current_span().await; - } - - PlaybackMessage::ToggleShuffle { span } => { - let _e = span.enter(); - debug!("toggling shuffle"); - let shuffle; - let repeat; - { - let Ok(mut queue) = self.queue.lock() else { - debug!("got queue lock"); - continue; - }; - debug!("got queue lock"); - repeat = queue.repeat; - if queue.shuffle { - queue.shuffle_off() - } else { - queue.shuffle_on() - } - shuffle = queue.shuffle; - } - debug!("queue lock released"); - let queue_update_tx = self.update_tx.clone(); - let update = StreamUpdate::Mods(QueueModifiers { shuffle, repeat }); - if let Err(err) = queue_update_tx.send(update) { - trace!("{:?}", err) - } - } - - PlaybackMessage::ToggleRepeat { span } => { - let _e = span.enter(); - debug!("toggling repeat"); - let shuffle; - let repeat; - { - let Ok(mut queue) = self.queue.lock() else { - debug!("got queue lock"); - continue; - }; - debug!("got queue lock"); - shuffle = queue.shuffle; - if queue.repeat { - queue.repeat = false - } else { - queue.repeat = true - } - repeat = queue.repeat; - } - debug!("queue lock released"); - let queue_update_tx = self.update_tx.clone(); - let update = StreamUpdate::Mods(QueueModifiers { shuffle, repeat }); - if let Err(err) = queue_update_tx.send(update) { - trace!("{:?}", err) - } - } - - PlaybackMessage::TogglePlay { span } => { - let _e = span.enter(); - debug!("toggling play"); - { - let state = { - let Ok(state) = self.state.lock() else { - debug!("got state lock"); - continue; - }; - *state - }; - debug!("got state lock"); - if state == PlayState::Playing { - if let Err(err) = self.player.pause().await { - error!("{:?}", err) - } - } else if let Err(err) = self.player.unpause().await { - error!("{:?}", err) - } - } - debug!("state lock released"); - } - - PlaybackMessage::Stop { span } => { - let _e = span.enter(); - debug!("stopping"); - if let Err(err) = self.player.stop().await { - error!("{:?}", err) - } - } - - PlaybackMessage::ChangeVolume { delta, span } => { - let _e = span.enter(); - debug!("changing volume"); - if let Ok(volume) = self.player.volume().await { - debug!("got volume {:?}", volume); - if let Err(err) = self.player.set_volume(volume + delta).await { - error!("{:?}", err) - }; - } - } - - PlaybackMessage::ToggleMute { span } => { - let _e = span.enter(); - debug!("toggling mute"); - // let muted = self.player.is_muted(); - // debug!("got muted {:?}", muted); - // self.player.set_mute(!muted); - } - - PlaybackMessage::Next { span } => { - let _e = span.enter(); - debug!("nexting"); - let track = { - let Ok(mut queue) = self.queue.lock() else { - debug!("got queue lock"); - continue; - }; - debug!("got queue lock"); - queue.next_track() - }; - debug!("released queue lock and got track {:?}", track); - - self.play_or_stop(track).in_current_span().await; - } - - PlaybackMessage::Prev { span } => { - let _e = span.enter(); - debug!("preving"); - let track = { - let Ok(mut queue) = self.queue.lock() else { - debug!("got queue lock"); - continue; - }; - debug!("got queue lock"); - queue.prev_track() - }; - debug!("released queue lock and got track {:?}", track); - self.play_or_stop(track).in_current_span().await; - } - - PlaybackMessage::StateChanged { state, span } => { - let _e = span.enter(); - debug!("state changed"); - - let play_state = { - let Ok(mut state_lock) = self.state.lock() else { - debug!("got state lock"); - continue; - }; - *state_lock = state; - state - }; - debug!("released state lock and got play state {:?}", play_state); - let active_track_tx = self.update_tx.clone(); - let update = StreamUpdate::PlayState(play_state as i32); - if let Err(err) = active_track_tx.send(update) { - trace!("{:?}", err) - }; - } - - PlaybackMessage::RestartTrack { span } => { - let _e = span.enter(); - debug!("restarting track"); - if let Err(err) = self.player.restart().await { - error!("{:?}", err) - } - } - - PlaybackMessage::VolumeChanged { volume, span } => { - let _e = span.enter(); - trace!("volume changed"); - let update_tx = self.update_tx.clone(); - let update = StreamUpdate::Volume(volume); - if let Err(err) = update_tx.send(update) { - trace!("{:?}", err) - } - } - - PlaybackMessage::MuteChanged { muted, span } => { - let _e = span.enter(); - trace!("mute changed"); - let update_tx = self.update_tx.clone(); - let update = StreamUpdate::Mute(muted); - if let Err(err) = update_tx.send(update) { - trace!("{:?}", err) - } - } - - PlaybackMessage::PostitionChanged { - duration, - position, - span, - } => { - let _e = span.enter(); - trace!("position changed"); - let update_tx = self.update_tx.clone(); - let update = StreamUpdate::Position(TrackPosition { duration, position }); - if let Err(err) = update_tx.send(update) { - trace!("{:?}", err) - } + PlaybackCommand::Remove { positions } => { + debug!(?positions, "removing tracks"); + let (track, was_last) = { + let Ok(mut queue) = self.queue.lock() else { + error!("queue lock poisoned"); + return; + }; + let was_last = queue.is_last_track(); + let track = queue.remove_tracks(&positions); + self.broadcast(StreamUpdate::Queue(queue.clone().into())); + (track, was_last) + }; + let state = { + let Ok(state) = self.state.lock() else { + error!("play state lock poisoned"); + return; + }; + *state + }; + if state == PlayState::Playing && track.is_some() { + // The playing track was removed: play its successor, or + // stop when it was the last one. + if was_last { + self.stop_player().await; + } else { + self.play(track).await; } } } - }); + + PlaybackCommand::Insert { position, uuids } => { + let all_tracks = self.resolve_tracks(uuids).await; + debug!(count = all_tracks.len(), position, "inserting into queue"); + let track = { + let Ok(mut queue) = self.queue.lock() else { + error!("queue lock poisoned"); + return; + }; + let track = queue.insert_tracks(position, &all_tracks); + self.broadcast(StreamUpdate::Queue(queue.clone().into())); + track + }; + self.play_if_some(track).await; + } + + PlaybackCommand::Clear { exclude_current } => { + debug!(exclude_current, "clearing queue"); + let should_stop = { + let Ok(mut queue) = self.queue.lock() else { + error!("queue lock poisoned"); + return; + }; + let should_stop = queue.clear(exclude_current); + self.broadcast(StreamUpdate::Queue(queue.clone().into())); + should_stop + }; + if should_stop { + self.stop_player().await; + } + } + + PlaybackCommand::SetCurrent { position } => { + debug!(position, "jumping to queue position"); + let track = { + let Ok(mut queue) = self.queue.lock() else { + error!("queue lock poisoned"); + return; + }; + queue.set_current_position(position); + queue.current_track() + }; + self.play(track).await; + } + + PlaybackCommand::ToggleShuffle => { + let (shuffle, repeat) = { + let Ok(mut queue) = self.queue.lock() else { + error!("queue lock poisoned"); + return; + }; + if queue.shuffle { + queue.shuffle_off() + } else { + queue.shuffle_on() + } + (queue.shuffle, queue.repeat) + }; + debug!(shuffle, "toggled shuffle"); + self.broadcast(StreamUpdate::Mods(QueueModifiers { shuffle, repeat })); + } + + PlaybackCommand::ToggleRepeat => { + let (shuffle, repeat) = { + let Ok(mut queue) = self.queue.lock() else { + error!("queue lock poisoned"); + return; + }; + queue.repeat = !queue.repeat; + (queue.shuffle, queue.repeat) + }; + debug!(repeat, "toggled repeat"); + self.broadcast(StreamUpdate::Mods(QueueModifiers { shuffle, repeat })); + } + + PlaybackCommand::TogglePlay => { + let state = { + let Ok(state) = self.state.lock() else { + error!("play state lock poisoned"); + return; + }; + *state + }; + debug!(?state, "toggling play"); + if state == PlayState::Playing { + if let Err(err) = self.player.pause().await { + warn!("pause failed: {err:?}"); + } + } else if let Err(err) = self.player.unpause().await { + warn!("unpause failed: {err:?}"); + } + } + + PlaybackCommand::Stop => { + debug!("stopping playback"); + self.stop_player().await; + } + + PlaybackCommand::ChangeVolume { delta } => { + match self.player.volume().await { + Ok(volume) => { + debug!(volume, delta, "changing volume"); + if let Err(err) = self.player.set_volume(volume + delta).await { + warn!("set_volume failed: {err:?}"); + } + } + Err(err) => warn!("could not read volume: {err:?}"), + }; + } + + PlaybackCommand::ToggleMute => { + // FIXME: implement mute in the player engine + debug!("toggle mute requested (not implemented)"); + } + + PlaybackCommand::Next => { + let track = { + let Ok(mut queue) = self.queue.lock() else { + error!("queue lock poisoned"); + return; + }; + queue.next_track() + }; + debug!( + track = track.as_ref().map(|t| t.uuid.as_str()), + "advancing to next track" + ); + self.play_or_stop(track).await; + } + + PlaybackCommand::Prev => { + let track = { + let Ok(mut queue) = self.queue.lock() else { + error!("queue lock poisoned"); + return; + }; + queue.prev_track() + }; + debug!( + track = track.as_ref().map(|t| t.uuid.as_str()), + "going back to previous track" + ); + self.play_or_stop(track).await; + } + + PlaybackCommand::StateChanged { state } => { + let play_state = { + let Ok(mut state_lock) = self.state.lock() else { + error!("play state lock poisoned"); + return; + }; + *state_lock = state; + state + }; + debug!(?play_state, "player state changed"); + self.broadcast(StreamUpdate::PlayState(play_state as i32)); + } + + PlaybackCommand::RestartTrack => { + debug!("restarting current track"); + if let Err(err) = self.player.restart().await { + warn!("restart failed: {err:?}"); + } + } + + PlaybackCommand::VolumeChanged { volume } => { + trace!(volume, "volume changed"); + self.broadcast(StreamUpdate::Volume(volume)); + } + + PlaybackCommand::MuteChanged { muted } => { + trace!(muted, "mute changed"); + self.broadcast(StreamUpdate::Mute(muted)); + } + + PlaybackCommand::PositionChanged { duration, position } => { + trace!(duration, position, "position changed"); + self.broadcast(StreamUpdate::Position(TrackPosition { duration, position })); + } + } + } + + /// Sends an update to all connected clients. Having no subscribers is + /// normal and not an error. + fn broadcast(&self, update: StreamUpdate) { + if let Err(err) = self.update_tx.send(update) { + trace!("no update stream subscribers: {err}"); + } + } + + /// Resolves a mixed list of track and node identifiers into tracks by + /// asking the provider orchestrator. + async fn resolve_tracks(&self, uuids: Vec) -> Vec { + let mut all_tracks = Vec::new(); + for uuid in uuids { + if is_track(&uuid) { + match self.get_track(&uuid).await { + Ok(track) => all_tracks.push(track), + Err(err) => warn!(uuid, "failed to resolve track: {err}"), + } + } else { + let tracks = self.flatten_node(&uuid).await; + all_tracks.extend(tracks); + } + } + trace!(count = all_tracks.len(), "resolved tracks"); + all_tracks } #[instrument(skip(self))] async fn flatten_node(&self, uuid: &str) -> Vec { - debug!("flattening node"); - let tx = self.provider_tx.clone(); let (result_tx, result_rx) = flume::bounded(1); - let span = debug_span!("prov-chan"); - let Ok(_) = tx - .send_async(ProviderMessage::FlattenNode { - uuid: uuid.to_string(), - result_tx, - span, - }) - .in_current_span() - .await - else { + let message = ProviderMessage::new(ProviderCommand::FlattenNode { + uuid: uuid.to_string(), + result_tx, + }); + if let Err(err) = self.provider_tx.send_async(message).await { + error!("provider channel closed: {err}"); return Vec::new(); - }; - let Ok(tracks) = result_rx.recv_async().in_current_span().await else { - return Vec::new(); - }; - tracks + } + match result_rx.recv_async().await { + Ok(tracks) => tracks, + Err(err) => { + error!("provider dropped flatten_node reply: {err}"); + Vec::new() + } + } } #[instrument(skip(self))] async fn get_track(&self, uuid: &str) -> Result { - debug!("getting track"); - let tx = self.provider_tx.clone(); let (result_tx, result_rx) = flume::bounded(1); - let span = tracing::trace_span!("prov-chan"); - tx.send_async(ProviderMessage::GetTrack { + let message = ProviderMessage::new(ProviderCommand::GetTrack { uuid: uuid.to_string(), result_tx, - span, - }) - .in_current_span() - .await - .map_err(|_| ProviderError::InternalError)?; + }); + self.provider_tx + .send_async(message) + .await + .map_err(|_| ProviderError::InternalError)?; result_rx .recv_async() - .in_current_span() .await .map_err(|_| ProviderError::InternalError)? } #[instrument(skip(self))] async fn get_urls_for_track(&self, uuid: &str) -> Result, ProviderError> { - debug!("getting urls for track"); - let tx = self.provider_tx.clone(); let (result_tx, result_rx) = flume::bounded(1); - let span = tracing::trace_span!("prov-chan"); - tx.send_async(ProviderMessage::GetTrackUrls { + let message = ProviderMessage::new(ProviderCommand::GetTrackUrls { uuid: uuid.to_string(), result_tx, - span, - }) - .in_current_span() - .await - .map_err(|_| ProviderError::InternalError)?; + }); + self.provider_tx + .send_async(message) + .await + .map_err(|_| ProviderError::InternalError)?; result_rx .recv_async() - .in_current_span() .await .map_err(|_| ProviderError::InternalError)? } - #[instrument(skip(self))] - async fn play_or_stop(&self, track: Option) { - debug!("play or stop"); - if let Some(track) = track { - let mut uuid = track.uuid.clone(); - let urls = loop { - match self.get_urls_for_track(&uuid).in_current_span().await { - Ok(urls) => break urls, - Err(err) => { - warn!("no urls found for track {:?}: {}", track.uuid, err); - uuid = { - let Ok(mut queue) = self.queue.lock() else { - debug!("got queue lock"); - continue; - }; - if let Some(track) = queue.next_track() { - track.uuid.clone() - } else { - return; - } - } - } - } - }; - { - let Ok(queue) = self.queue.lock() else { - error!("poisend queue lock"); - return; - }; - let queue_update_tx = self.update_tx.clone(); - let track = queue.current_track(); - let update = StreamUpdate::QueueTrack(QueueTrack { - queue_position: queue.current_position() as u32, - track, - }); - if let Err(err) = queue_update_tx.send(update) { - trace!("{:?}", err) - } - } - if let Err(err) = self.player.play(&urls[0]).await { - error!("{:?}", err) - }; - } else if let Err(err) = self.player.stop().await { - error!("{:?}", err) + async fn stop_player(&self) { + if let Err(err) = self.player.stop().await { + debug!("stop had no effect: {err:?}"); } } - #[instrument(skip(self))] + /// Plays the given track if there is one, otherwise stops the player. + #[instrument(skip(self, track), fields(track = track.as_ref().map(|t| t.uuid.as_str())))] + async fn play_or_stop(&self, track: Option) { + if track.is_some() { + self.play(track).await; + } else { + self.stop_player().await; + } + } + + /// Plays the given track if there is one; does nothing otherwise. + async fn play_if_some(&self, track: Option) { + if track.is_some() { + self.play(track).await; + } + } + + /// Starts playback of the given track. When fetching stream URLs fails + /// the failing track is skipped and playback continues with the next + /// track in the queue. + #[instrument(skip(self, track), fields(track = track.as_ref().map(|t| t.uuid.as_str())))] async fn play(&self, track: Option) { - debug!("play"); - if let Some(track) = track { - let mut uuid = track.uuid.clone(); - let urls = loop { - match self.get_urls_for_track(&uuid).in_current_span().await { - Ok(urls) => break urls, - Err(err) => { - warn!("no urls found for track {:?}: {}", track.uuid, err); - uuid = { - let Ok(mut queue) = self.queue.lock() else { - debug!("poisend queue lock"); - return; - }; - if let Some(track) = queue.next_track() { - track.uuid.clone() - } else { - return; - } - } - } - } - }; - { - let Ok(queue) = self.queue.lock() else { - error!("poisend queue lock"); + let Some(track) = track else { + debug!("nothing to play"); + return; + }; + let mut uuid = track.uuid.clone(); + let urls = loop { + match self.get_urls_for_track(&uuid).await { + Ok(urls) if !urls.is_empty() => break urls, + Ok(_) => warn!(uuid, "provider returned no stream urls, skipping track"), + Err(err) => warn!(uuid, "failed to fetch stream urls ({err}), skipping track"), + } + let next = { + let Ok(mut queue) = self.queue.lock() else { + error!("queue lock poisoned"); return; }; - let queue_update_tx = self.update_tx.clone(); - let track = queue.current_track(); - let update = StreamUpdate::QueueTrack(QueueTrack { - queue_position: queue.current_position() as u32, - track, - }); - if let Err(err) = queue_update_tx.send(update) { - trace!("{:?}", err) + queue.next_track() + }; + match next { + Some(next_track) => uuid = next_track.uuid.clone(), + None => { + error!("no playable track left in queue, stopping"); + self.stop_player().await; + return; } } - if let Err(err) = self.player.play(&urls[0]).await { - error!("{:?}", err) - } + }; + { + let Ok(queue) = self.queue.lock() else { + error!("queue lock poisoned"); + return; + }; + self.broadcast(StreamUpdate::QueueTrack(QueueTrack { + queue_position: queue.current_position() as u32, + track: queue.current_track(), + })); + } + debug!(url_count = urls.len(), "starting player"); + if let Err(err) = self.player.play(&urls[0]).await { + error!("player failed to start track: {err:?}"); } } } diff --git a/crabidy-server/src/provider.rs b/crabidy-server/src/provider.rs index 3992388..fe25115 100644 --- a/crabidy-server/src/provider.rs +++ b/crabidy-server/src/provider.rs @@ -1,88 +1,80 @@ -use crate::ProviderMessage; +use crate::{ProviderCommand, ProviderMessage}; use async_trait::async_trait; use crabidy_core::{ proto::crabidy::{LibraryNode, LibraryNodeChild, Track}, ProviderClient, ProviderError, }; use std::{fs, path::PathBuf, sync::Arc}; -use tracing::{debug, error, instrument, warn, Instrument}; +use tracing::{debug, debug_span, error, instrument, warn, Instrument}; #[derive(Debug)] pub struct ProviderOrchestrator { pub provider_tx: flume::Sender, provider_rx: flume::Receiver, - // known_tracks: RwLock>, - // known_nodes: RwLock>, tidal_client: Arc, } impl ProviderOrchestrator { pub fn run(self) { tokio::spawn(async move { - while let Ok(msg) = self.provider_rx.recv_async().await { - match msg { - ProviderMessage::GetLibraryNode { - uuid, - result_tx, - span, - } => { - let _e = span.enter(); - let result = self.get_lib_node(&uuid).in_current_span().await; - if let Err(err) = result_tx.send_async(result).in_current_span().await { - error!("failed to send result: {}", err); - } - } - ProviderMessage::GetTrack { - uuid, - result_tx, - span, - } => { - let _e = span.enter(); - let result = self.get_metadata_for_track(&uuid).in_current_span().await; - if let Err(err) = result_tx.send_async(result).in_current_span().await { - error!("failed to send result: {}", err); - } - } - ProviderMessage::GetTrackUrls { - uuid, - result_tx, - span, - } => { - let _e = span.enter(); - let result = self.get_urls_for_track(&uuid).in_current_span().await; - if let Err(err) = result_tx.send_async(result).in_current_span().await { - error!("failed to send result: {}", err); - } - } - ProviderMessage::FlattenNode { - uuid, - result_tx, - span, - } => { - let _e = span.enter(); - let result = self.flatten_node(&uuid).in_current_span().await; - if let Err(err) = result_tx.send_async(result).in_current_span().await { - error!("failed to send result: {}", err); - } - } - } + while let Ok(ProviderMessage { span, command }) = self.provider_rx.recv_async().await { + let handler_span = + debug_span!(parent: &span, "provider_command", command = command.name()); + self.handle_command(command).instrument(handler_span).await; } + warn!("provider message channel closed, loop exiting"); }); } + + async fn handle_command(&self, command: ProviderCommand) { + match command { + ProviderCommand::GetLibraryNode { uuid, result_tx } => { + let result = self.get_lib_node(&uuid).await; + if let Err(err) = result_tx.send_async(result).await { + error!("failed to send get_library_node result: {err}"); + } + } + ProviderCommand::GetTrack { uuid, result_tx } => { + let result = self.get_metadata_for_track(&uuid).await; + if let Err(err) = result_tx.send_async(result).await { + error!("failed to send get_track result: {err}"); + } + } + ProviderCommand::GetTrackUrls { uuid, result_tx } => { + let result = self.get_urls_for_track(&uuid).await; + if let Err(err) = result_tx.send_async(result).await { + error!("failed to send get_track_urls result: {err}"); + } + } + ProviderCommand::FlattenNode { uuid, result_tx } => { + let result = self.flatten_node(&uuid).await; + if let Err(err) = result_tx.send_async(result).await { + error!("failed to send flatten_node result: {err}"); + } + } + } + } + + /// Collects all tracks reachable from the given node by walking its + /// queueable descendants. #[instrument(skip(self))] async fn flatten_node(&self, node_uuid: &str) -> Vec { - let mut tracks = Vec::with_capacity(1000); - let mut nodes_to_go = Vec::with_capacity(100); - nodes_to_go.push(node_uuid.to_string()); + let mut tracks = Vec::new(); + let mut nodes_to_go = vec![node_uuid.to_string()]; while let Some(node_uuid) = nodes_to_go.pop() { - let Ok(node) = self.get_lib_node(&node_uuid).in_current_span().await else { - continue; + let node = match self.get_lib_node(&node_uuid).await { + Ok(node) => node, + Err(err) => { + warn!(node = node_uuid, "skipping unreadable node: {err}"); + continue; + } }; if node.is_queable { tracks.extend(node.tracks); nodes_to_go.extend(node.children.into_iter().map(|c| c.uuid)) } } + debug!(count = tracks.len(), "flattened node into tracks"); tracks } } @@ -95,29 +87,25 @@ impl ProviderClient for ProviderOrchestrator { .map(|d| d.join("crabidy")) .unwrap_or(PathBuf::from("/tmp")); let dir_exists = tokio::fs::try_exists(&config_dir) - .in_current_span() .await .map_err(|e| ProviderError::Config(e.to_string()))?; if !dir_exists { tokio::fs::create_dir(&config_dir) - .in_current_span() .await .map_err(|e| ProviderError::Config(e.to_string()))?; } let config_file = config_dir.join("tidaly.toml"); - let raw_toml_settings = fs::read_to_string(&config_file).unwrap_or("".to_owned()); - let tidal_client = Arc::new( - tidaldy::Client::init(&raw_toml_settings) - .in_current_span() - .await - .expect("Failed to init Tidal clienta"), - ); + debug!(config_file = %config_file.display(), "loading tidal config"); + let raw_toml_settings = fs::read_to_string(&config_file).unwrap_or_default(); + let tidal_client = Arc::new(tidaldy::Client::init(&raw_toml_settings).await.map_err( + |err| { + error!("failed to init tidal client: {err}"); + err + }, + )?); let new_toml_config = tidal_client.settings(); - if let Err(err) = tokio::fs::write(&config_file, new_toml_config) - .in_current_span() - .await - { - error!("Failed to write config file: {}", err); + if let Err(err) = tokio::fs::write(&config_file, new_toml_config).await { + error!("failed to write tidal config file: {err}"); }; let (provider_tx, provider_rx) = flume::bounded(100); Ok(Self { @@ -126,46 +114,38 @@ impl ProviderClient for ProviderOrchestrator { tidal_client, }) } - #[instrument(skip(self))] + fn settings(&self) -> String { - "".to_owned() + String::new() } + #[instrument(skip(self))] async fn get_urls_for_track(&self, track_uuid: &str) -> Result, ProviderError> { - debug!("get_urls_for_track"); - self.tidal_client - .get_urls_for_track(track_uuid) - .in_current_span() - .await + self.tidal_client.get_urls_for_track(track_uuid).await } + #[instrument(skip(self))] async fn get_metadata_for_track(&self, track_uuid: &str) -> Result { - debug!("get_metadata_for_track"); - self.tidal_client - .get_metadata_for_track(track_uuid) - .in_current_span() - .await + self.tidal_client.get_metadata_for_track(track_uuid).await } - #[instrument(skip(self))] + fn get_lib_root(&self) -> LibraryNode { - debug!("get_lib_root in provider manager"); let mut root_node = LibraryNode::new(); let child = LibraryNodeChild::new("node:tidal".to_owned(), "tidal".to_owned(), false); root_node.children.push(child); root_node } + #[instrument(skip(self))] async fn get_lib_node(&self, uuid: &str) -> Result { - debug!("get_lib_node in provider manager"); if uuid == "node:/" { - debug!("get global root"); + debug!("serving global library root"); return Ok(self.get_lib_root()); } if uuid == "node:tidal" { - debug!("get tidal root"); + debug!("serving tidal library root"); return Ok(self.tidal_client.get_lib_root()); } - debug!("tidal node"); - self.tidal_client.get_lib_node(uuid).in_current_span().await + self.tidal_client.get_lib_node(uuid).await } } diff --git a/crabidy-server/src/rpc.rs b/crabidy-server/src/rpc.rs index 61bb81d..a4de109 100644 --- a/crabidy-server/src/rpc.rs +++ b/crabidy-server/src/rpc.rs @@ -1,4 +1,4 @@ -use crate::{PlaybackMessage, ProviderMessage}; +use crate::{PlaybackCommand, PlaybackMessage, ProviderCommand, ProviderMessage}; use crabidy_core::proto::crabidy::{ crabidy_service_server::CrabidyService, get_update_stream_response::Update as StreamUpdate, AppendRequest, AppendResponse, ChangeVolumeRequest, ChangeVolumeResponse, ClearQueueRequest, @@ -10,11 +10,10 @@ use crabidy_core::proto::crabidy::{ StopResponse, ToggleMuteRequest, ToggleMuteResponse, TogglePlayRequest, TogglePlayResponse, ToggleRepeatRequest, ToggleRepeatResponse, ToggleShuffleRequest, ToggleShuffleResponse, }; -use futures::TryStreamExt; use std::pin::Pin; use tokio_stream::StreamExt; use tonic::{Request, Response, Status}; -use tracing::{debug, debug_span, error, instrument, trace, Instrument, Span}; +use tracing::{debug, error, instrument, trace}; #[derive(Debug)] pub struct RpcService { @@ -25,16 +24,29 @@ pub struct RpcService { impl RpcService { pub fn new( - update_rx: tokio::sync::broadcast::Sender, + update_tx: tokio::sync::broadcast::Sender, playback_tx: flume::Sender, provider_tx: flume::Sender, ) -> Self { Self { - update_tx: update_rx, + update_tx, playback_tx, provider_tx, } } + + /// Sends a command to the playback loop, mapping channel failure to an + /// internal error status. + async fn send_playback(&self, command: PlaybackCommand) -> Result<(), Status> { + let name = command.name(); + self.playback_tx + .send_async(PlaybackMessage::new(command)) + .await + .map_err(|err| { + error!(command = name, "playback channel closed: {err}"); + Status::internal("playback loop unavailable") + }) + } } #[tonic::async_trait] @@ -44,26 +56,14 @@ impl CrabidyService for RpcService { #[instrument(skip(self, _request))] async fn init(&self, _request: Request) -> Result, Status> { - debug!("Received init request"); - let playback_tx = self.playback_tx.clone(); + debug!("received init request"); let (result_tx, result_rx) = flume::bounded(1); - let span = debug_span!("play-chan"); - if let Err(err) = playback_tx - .send_async(PlaybackMessage::Init { result_tx, span }) - .in_current_span() - .await - { - error!("{:?}", err); - return Err(Status::internal("Sending Init via internal channel failed")); - } - let response = result_rx - .recv_async() - .in_current_span() - .await - .map_err(|e| { - error!("{:?}", e); - Status::internal("Failed to receive response from provider channel") - })?; + self.send_playback(PlaybackCommand::Init { result_tx }) + .await?; + let response = result_rx.recv_async().await.map_err(|err| { + error!("no reply from playback loop: {err}"); + Status::internal("playback loop did not reply") + })?; Ok(Response::new(response)) } @@ -73,32 +73,27 @@ impl CrabidyService for RpcService { request: Request, ) -> Result, Status> { let uuid = request.into_inner().uuid; - Span::current().record("uuid", &uuid); - debug!("Received get_library_node request"); - let provider_tx = self.provider_tx.clone(); + tracing::Span::current().record("uuid", uuid.as_str()); + debug!("received get_library_node request"); let (result_tx, result_rx) = flume::bounded(1); - let span = debug_span!("prov-chan"); - provider_tx - .send_async(ProviderMessage::GetLibraryNode { + self.provider_tx + .send_async(ProviderMessage::new(ProviderCommand::GetLibraryNode { uuid, result_tx, - span, - }) - .in_current_span() + })) .await - .map_err(|_| Status::internal("Failed to send request via channel"))?; - let result = result_rx - .recv_async() - .in_current_span() - .await - .map_err(|e| { - error!("{:?}", e); - Status::internal("Failed to receive response from provider channel") + .map_err(|err| { + error!("provider channel closed: {err}"); + Status::internal("provider unavailable") })?; + let result = result_rx.recv_async().await.map_err(|err| { + error!("no reply from provider: {err}"); + Status::internal("provider did not reply") + })?; match result { Ok(node) => Ok(Response::new(GetLibraryNodeResponse { node: Some(node) })), Err(err) => { - error!("{:?}", err); + error!("get_library_node failed: {err}"); Err(Status::internal(err.to_string())) } } @@ -107,204 +102,137 @@ impl CrabidyService for RpcService { #[instrument(skip(self, request), fields(uuids))] async fn queue( &self, - request: tonic::Request, - ) -> std::result::Result, tonic::Status> { - let uuids = request.into_inner().uuids.clone(); - Span::current().record("uuids", format!("{:?}", uuids)); - debug!("Received queue request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - playback_tx - .send_async(PlaybackMessage::Queue { uuids, span }) - .in_current_span() - .await - .map_err(|_| Status::internal("Failed to send request via channel"))?; - - let reply = QueueResponse {}; - Ok(Response::new(reply)) + request: Request, + ) -> Result, Status> { + let uuids = request.into_inner().uuids; + tracing::Span::current().record("uuids", format!("{uuids:?}")); + debug!("received queue request"); + self.send_playback(PlaybackCommand::Queue { uuids }).await?; + Ok(Response::new(QueueResponse {})) } #[instrument(skip(self, request), fields(uuids))] async fn replace( &self, - request: tonic::Request, - ) -> std::result::Result, tonic::Status> { - let uuids = request.into_inner().uuids.clone(); - Span::current().record("uuids", format!("{:?}", uuids)); - debug!("Received replace request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - playback_tx - .send_async(PlaybackMessage::Replace { uuids, span }) - .in_current_span() - .await - .map_err(|_| Status::internal("Failed to send request via channel"))?; - let reply = ReplaceResponse {}; - Ok(Response::new(reply)) + request: Request, + ) -> Result, Status> { + let uuids = request.into_inner().uuids; + tracing::Span::current().record("uuids", format!("{uuids:?}")); + debug!("received replace request"); + self.send_playback(PlaybackCommand::Replace { uuids }) + .await?; + Ok(Response::new(ReplaceResponse {})) } #[instrument(skip(self, request), fields(uuids))] async fn append( &self, - request: tonic::Request, - ) -> std::result::Result, tonic::Status> { - let uuids = request.into_inner().uuids.clone(); - Span::current().record("uuids", format!("{:?}", uuids)); - debug!("Received append request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - playback_tx - .send_async(PlaybackMessage::Append { uuids, span }) - .in_current_span() - .await - .map_err(|_| Status::internal("Failed to send request via channel"))?; - let reply = AppendResponse {}; - Ok(Response::new(reply)) + request: Request, + ) -> Result, Status> { + let uuids = request.into_inner().uuids; + tracing::Span::current().record("uuids", format!("{uuids:?}")); + debug!("received append request"); + self.send_playback(PlaybackCommand::Append { uuids }) + .await?; + Ok(Response::new(AppendResponse {})) } #[instrument(skip(self, request), fields(positions))] async fn remove( &self, - request: tonic::Request, - ) -> std::result::Result, tonic::Status> { + request: Request, + ) -> Result, Status> { let positions = request.into_inner().positions; - Span::current().record("positions", format!("{:?}", positions)); - debug!("Received remove request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - playback_tx - .send_async(PlaybackMessage::Remove { positions, span }) - .in_current_span() - .await - .map_err(|_| Status::internal("Failed to send request via channel"))?; - let reply = RemoveResponse {}; - Ok(Response::new(reply)) + tracing::Span::current().record("positions", format!("{positions:?}")); + debug!("received remove request"); + self.send_playback(PlaybackCommand::Remove { positions }) + .await?; + Ok(Response::new(RemoveResponse {})) } #[instrument(skip(self, request), fields(uuids, position))] async fn insert( &self, - request: tonic::Request, - ) -> std::result::Result, tonic::Status> { + request: Request, + ) -> Result, Status> { let req = request.into_inner(); - let uuids = req.uuids.clone(); - let position = req.position; - Span::current().record("uuids", format!("{:?}", uuids)); - Span::current().record("position", position); - debug!("Received insert request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - playback_tx - .send_async(PlaybackMessage::Insert { - position: req.position, - uuids, - span, - }) - .in_current_span() - .await - .map_err(|_| Status::internal("Failed to send request via channel"))?; - let reply = InsertResponse {}; - Ok(Response::new(reply)) + tracing::Span::current().record("uuids", format!("{:?}", req.uuids)); + tracing::Span::current().record("position", req.position); + debug!("received insert request"); + self.send_playback(PlaybackCommand::Insert { + position: req.position, + uuids: req.uuids, + }) + .await?; + Ok(Response::new(InsertResponse {})) } #[instrument(skip(self, request), fields(exclude_current))] async fn clear_queue( &self, - request: tonic::Request, - ) -> std::result::Result, tonic::Status> { + request: Request, + ) -> Result, Status> { let exclude_current = request.into_inner().exclude_current; - Span::current().record("exclude_current", exclude_current); - debug!("Received clear_queue request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - playback_tx - .send_async(PlaybackMessage::Clear { - exclude_current, - span, - }) - .in_current_span() - .await - .map_err(|_| Status::internal("Failed to send request via channel"))?; - let reply = ClearQueueResponse {}; - Ok(Response::new(reply)) + tracing::Span::current().record("exclude_current", exclude_current); + debug!("received clear_queue request"); + self.send_playback(PlaybackCommand::Clear { exclude_current }) + .await?; + Ok(Response::new(ClearQueueResponse {})) } #[instrument(skip(self, request), fields(position))] async fn set_current( &self, - request: tonic::Request, - ) -> std::result::Result, tonic::Status> { + request: Request, + ) -> Result, Status> { let position = request.into_inner().position; - Span::current().record("position", position); - debug!("Received set_current request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - playback_tx - .send_async(PlaybackMessage::SetCurrent { position, span }) - .in_current_span() - .await - .map_err(|_| Status::internal("Failed to send request via channel"))?; - let reply = SetCurrentResponse {}; - Ok(Response::new(reply)) + tracing::Span::current().record("position", position); + debug!("received set_current request"); + self.send_playback(PlaybackCommand::SetCurrent { position }) + .await?; + Ok(Response::new(SetCurrentResponse {})) } #[instrument(skip(self, _request))] async fn toggle_shuffle( &self, - _request: tonic::Request, - ) -> std::result::Result, tonic::Status> { - debug!("Received toggle_shuffle request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - if let Err(err) = playback_tx - .send_async(PlaybackMessage::ToggleShuffle { span }) - .in_current_span() - .await - { - error!("Failed to send request via channel: {}", err); - } - let reply = ToggleShuffleResponse {}; - Ok(Response::new(reply)) + _request: Request, + ) -> Result, Status> { + debug!("received toggle_shuffle request"); + self.send_playback(PlaybackCommand::ToggleShuffle).await?; + Ok(Response::new(ToggleShuffleResponse {})) } #[instrument(skip(self, _request))] async fn toggle_repeat( &self, - _request: tonic::Request, - ) -> std::result::Result, tonic::Status> { - debug!("Received toggle_repeat request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - if let Err(err) = playback_tx - .send_async(PlaybackMessage::ToggleRepeat { span }) - .in_current_span() - .await - { - error!("Failed to send request via channel: {}", err); - } - let reply = ToggleRepeatResponse {}; - Ok(Response::new(reply)) + _request: Request, + ) -> Result, Status> { + debug!("received toggle_repeat request"); + self.send_playback(PlaybackCommand::ToggleRepeat).await?; + Ok(Response::new(ToggleRepeatResponse {})) } #[instrument(skip(self, _request))] async fn get_update_stream( &self, - _request: tonic::Request, - ) -> std::result::Result, tonic::Status> { - debug!("Received get_update_stream request"); + _request: Request, + ) -> Result, Status> { + debug!("received get_update_stream request, subscribing client"); let update_rx = self.update_tx.subscribe(); let update_stream = tokio_stream::wrappers::BroadcastStream::new(update_rx); - let output_stream = update_stream.into_stream().map(|update_result| { - trace!("Got update: {:?}", update_result); + let output_stream = update_stream.map(|update_result| { + trace!(?update_result, "forwarding update"); match update_result { Ok(update) => Ok(GetUpdateStreamResponse { update: Some(update), }), - Err(_) => Err(tonic::Status::new( - tonic::Code::Unknown, - "Internal channel error", - )), + Err(err) => { + // The client lagged too far behind the broadcast channel. + error!("update stream lagged: {err}"); + Err(Status::data_loss("update stream lagged")) + } } }); @@ -314,146 +242,73 @@ impl CrabidyService for RpcService { #[instrument(skip(self, _request))] async fn save_queue( &self, - _request: tonic::Request, - ) -> std::result::Result, tonic::Status> { - debug!("Received save_queue request"); - let reply = SaveQueueResponse {}; - Ok(Response::new(reply)) + _request: Request, + ) -> Result, Status> { + debug!("received save_queue request (not implemented)"); + Ok(Response::new(SaveQueueResponse {})) } - /// Playback #[instrument(skip(self, _request))] async fn toggle_play( &self, - _request: tonic::Request, - ) -> std::result::Result, tonic::Status> { - debug!("Received toggle_play request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - if let Err(err) = playback_tx - .send_async(PlaybackMessage::TogglePlay { span }) - .in_current_span() - .await - { - error!("Failed to send request via channel: {}", err); - } - let reply = TogglePlayResponse {}; - Ok(Response::new(reply)) + _request: Request, + ) -> Result, Status> { + debug!("received toggle_play request"); + self.send_playback(PlaybackCommand::TogglePlay).await?; + Ok(Response::new(TogglePlayResponse {})) } #[instrument(skip(self, _request))] - async fn stop( - &self, - _request: tonic::Request, - ) -> std::result::Result, tonic::Status> { - debug!("Received stop request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - if let Err(err) = playback_tx - .send_async(PlaybackMessage::Stop { span }) - .in_current_span() - .await - { - error!("Failed to send request via channel: {}", err); - } - let reply = StopResponse {}; - Ok(Response::new(reply)) + async fn stop(&self, _request: Request) -> Result, Status> { + debug!("received stop request"); + self.send_playback(PlaybackCommand::Stop).await?; + Ok(Response::new(StopResponse {})) } #[instrument(skip(self, request), fields(delta))] async fn change_volume( &self, - request: tonic::Request, - ) -> std::result::Result, tonic::Status> { + request: Request, + ) -> Result, Status> { let delta = request.into_inner().delta; - Span::current().record("delta", delta); - debug!("Received change_volume request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - if let Err(err) = playback_tx - .send_async(PlaybackMessage::ChangeVolume { delta, span }) - .in_current_span() - .await - { - error!("Failed to send request via channel: {}", err); - } - let reply = ChangeVolumeResponse {}; - Ok(Response::new(reply)) + tracing::Span::current().record("delta", delta); + debug!("received change_volume request"); + self.send_playback(PlaybackCommand::ChangeVolume { delta }) + .await?; + Ok(Response::new(ChangeVolumeResponse {})) } #[instrument(skip(self, _request))] async fn toggle_mute( &self, - _request: tonic::Request, - ) -> std::result::Result, tonic::Status> { - debug!("Received toggle_mute request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - if let Err(err) = playback_tx - .send_async(PlaybackMessage::ToggleMute { span }) - .in_current_span() - .await - { - error!("Failed to send request via channel: {}", err); - } - let reply = ToggleMuteResponse {}; - Ok(Response::new(reply)) + _request: Request, + ) -> Result, Status> { + debug!("received toggle_mute request"); + self.send_playback(PlaybackCommand::ToggleMute).await?; + Ok(Response::new(ToggleMuteResponse {})) } #[instrument(skip(self, _request))] - async fn next( - &self, - _request: tonic::Request, - ) -> std::result::Result, tonic::Status> { - debug!("Received next request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - if let Err(err) = playback_tx - .send_async(PlaybackMessage::Next { span }) - .in_current_span() - .await - { - error!("Failed to send request via channel: {}", err); - } - let reply = NextResponse {}; - Ok(Response::new(reply)) + async fn next(&self, _request: Request) -> Result, Status> { + debug!("received next request"); + self.send_playback(PlaybackCommand::Next).await?; + Ok(Response::new(NextResponse {})) } #[instrument(skip(self, _request))] - async fn prev( - &self, - _request: tonic::Request, - ) -> std::result::Result, tonic::Status> { - debug!("Received prev request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - if let Err(err) = playback_tx - .send_async(PlaybackMessage::Prev { span }) - .in_current_span() - .await - { - error!("Failed to send request via channel: {}", err); - } - let reply = PrevResponse {}; - Ok(Response::new(reply)) + async fn prev(&self, _request: Request) -> Result, Status> { + debug!("received prev request"); + self.send_playback(PlaybackCommand::Prev).await?; + Ok(Response::new(PrevResponse {})) } #[instrument(skip(self, _request))] async fn restart_track( &self, - _request: tonic::Request, - ) -> std::result::Result, tonic::Status> { - debug!("Received restart_track request"); - let playback_tx = self.playback_tx.clone(); - let span = debug_span!("play-chan"); - if let Err(err) = playback_tx - .send_async(PlaybackMessage::RestartTrack { span }) - .in_current_span() - .await - { - error!("Failed to send request via channel: {}", err); - } - let reply = RestartTrackResponse {}; - Ok(Response::new(reply)) + _request: Request, + ) -> Result, Status> { + debug!("received restart_track request"); + self.send_playback(PlaybackCommand::RestartTrack).await?; + Ok(Response::new(RestartTrackResponse {})) } } diff --git a/tidaldy/src/lib.rs b/tidaldy/src/lib.rs index 42aaedd..4dd49c9 100644 --- a/tidaldy/src/lib.rs +++ b/tidaldy/src/lib.rs @@ -3,7 +3,7 @@ use reqwest::Client as HttpClient; use serde::de::DeserializeOwned; use tokio::time::{sleep, Duration, Instant}; -use tracing::{debug, error, info, instrument}; +use tracing::{debug, error, info, instrument, trace, warn}; pub mod config; pub mod models; use async_trait::async_trait; @@ -22,12 +22,8 @@ impl crabidy_core::ProviderClient for Client { let settings: config::Settings = if let Ok(settings) = toml::from_str(raw_toml_settings) { settings } else { - let settings = config::Settings::default(); - println!( - "could not parse toml settings: {:#?} using default settings instead: {:#?}", - raw_toml_settings, settings - ); - settings + warn!("could not parse toml settings, using defaults"); + config::Settings::default() }; let mut client = Self::new(settings)?; @@ -48,16 +44,21 @@ impl crabidy_core::ProviderClient for Client { &self, track_uuid: &str, ) -> Result, crabidy_core::ProviderError> { - debug!("get_urls_for_track {}", track_uuid); let (_, track_uuid, _) = split_uuid(track_uuid); - let Ok(playback) = self.get_track_playback(&track_uuid).await else { - return Err(crabidy_core::ProviderError::FetchError); - }; - debug!("playback {:?}", playback); - let Ok(manifest) = playback.get_manifest() else { - return Err(crabidy_core::ProviderError::FetchError); - }; - debug!("manifest {:?}", manifest); + let playback = self.get_track_playback(&track_uuid).await.map_err(|err| { + warn!(track = track_uuid, "failed to fetch playback info: {err}"); + crabidy_core::ProviderError::FetchError + })?; + trace!(?playback, "got playback info"); + let manifest = playback.get_manifest().map_err(|err| { + warn!(track = track_uuid, "failed to decode manifest: {err}"); + crabidy_core::ProviderError::FetchError + })?; + debug!( + track = track_uuid, + urls = manifest.urls.len(), + "resolved stream urls" + ); Ok(manifest.urls) } @@ -66,10 +67,10 @@ impl crabidy_core::ProviderClient for Client { &self, track_uuid: &str, ) -> Result { - debug!("get_metadata_for_track {}", track_uuid); - let Ok(track) = self.get_track(track_uuid).await else { - return Err(crabidy_core::ProviderError::FetchError); - }; + let track = self.get_track(track_uuid).await.map_err(|err| { + warn!(track = track_uuid, "failed to fetch track metadata: {err}"); + crabidy_core::ProviderError::FetchError + })?; Ok(track.into()) } @@ -107,9 +108,8 @@ impl crabidy_core::ProviderClient for Client { let Some(user_id) = self.settings.login.user_id.clone() else { return Err(crabidy_core::ProviderError::UnknownUser); }; - debug!("get_lib_node in tidaldy{}", uuid); let (_kind, module, uuid) = split_uuid(uuid); - error!("module:{},uuid: {}", module, uuid); + debug!(module, uuid, "resolving library node"); let node = match module.as_str() { "userplaylists" => { let mut node = crabidy_core::proto::crabidy::LibraryNode { @@ -166,7 +166,6 @@ impl crabidy_core::ProviderClient for Client { node } "artist" => { - info!("artist"); let mut node: crabidy_core::proto::crabidy::LibraryNode = self.get_artist(&uuid).await?.into(); let children: Vec = self @@ -371,7 +370,7 @@ impl Client { error!("{:?}", e); e })?; - println!("{:?}", response); + debug!(?response, "explorer response"); Ok(()) } @@ -504,7 +503,13 @@ impl Client { pub async fn login_web(&mut self) -> Result<(), ClientError> { let code_response = self.get_device_code().await?; let now = Instant::now(); + // The verification link must reach the user even without a log + // subscriber configured. println!("https://{}", code_response.verification_uri_complete); + info!( + "waiting for device login at https://{}", + code_response.verification_uri_complete + ); while now.elapsed().as_secs() <= code_response.expires_in { let login = self.check_auth_status(&code_response.device_code).await; if login.is_err() { @@ -520,9 +525,10 @@ impl Client { self.settings.login.expires_after = Some(login_results.expires_in + timestamp); self.settings.login.user_id = Some(login_results.user.user_id.to_string()); self.settings.login.country_code = Some(login_results.user.country_code); + info!("device login succeeded"); return Ok(()); } - println!("login attempt expired"); + warn!("device login attempt expired"); Err(ClientError::ConnectionError) }