Overhaul tracing: fix span misattribution, broaden coverage

The old pattern passed a Span in every channel message and entered it
with a guard that was held across await points, which misattributed
events from interleaved tasks. Messages are now {span, command} pairs:
the span is captured automatically at send time (Span::current) and the
consumer instruments the whole handler future with a child span
(playback_command/provider_command with a command name field), so events
are attributed correctly across the queue boundary and all the manual
in_current_span() plumbing is gone.

Also:
- server: EnvFilter with RUST_LOG support (default: own crates at
  debug, rest at info), log-crate bridge for symphonia/cpal
- cbd-tui: logs to a file under the state dir (the terminal belongs to
  the TUI), EnvFilter, no more println into the alternate screen
- tidaldy: fix misused levels (error->debug), structured fields,
  payload dumps moved to trace, login flow at info/warn
- no panic on missing notification daemon in the TUI
- no panic on backwards clock steps in QueueManager
- provider init errors propagate instead of expect()

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
AI User 2026-07-19 21:48:00 +02:00
parent 6eb5a87b15
commit 6d8bc7f166
11 changed files with 877 additions and 1119 deletions

4
Cargo.lock generated
View File

@ -454,6 +454,7 @@ version = "0.1.0"
dependencies = [ dependencies = [
"crabidy-core", "crabidy-core",
"crossterm", "crossterm",
"dirs",
"flume", "flume",
"notify-rust", "notify-rust",
"ratatui", "ratatui",
@ -461,6 +462,9 @@ dependencies = [
"tokio", "tokio",
"tokio-stream", "tokio-stream",
"tonic", "tonic",
"tracing",
"tracing-appender",
"tracing-subscriber",
] ]
[[package]] [[package]]

View File

@ -6,6 +6,7 @@ edition.workspace = true
[dependencies] [dependencies]
crabidy-core.workspace = true crabidy-core.workspace = true
crossterm.workspace = true crossterm.workspace = true
dirs.workspace = true
flume.workspace = true flume.workspace = true
notify-rust.workspace = true notify-rust.workspace = true
ratatui.workspace = true ratatui.workspace = true
@ -13,3 +14,6 @@ serde.workspace = true
tokio = { workspace = true, features = ["full"] } tokio = { workspace = true, features = ["full"] }
tokio-stream.workspace = true tokio-stream.workspace = true
tonic.workspace = true tonic.workspace = true
tracing.workspace = true
tracing-appender.workspace = true
tracing-subscriber.workspace = true

View File

@ -56,11 +56,14 @@ impl NowPlaying {
} else { } else {
format!("{} by {}", track.title, track.artist,) 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") .summary("Now playing")
.body(&body) .body(&body)
.show() .show()
.unwrap(); {
tracing::debug!("could not show desktop notification: {err}");
}
} }
self.track = active; self.track = active;
} }

View File

@ -27,11 +27,45 @@ use tokio_stream::StreamExt;
use app::{App, MessageFromUi, MessageToUi, StatefulList, UiFocus}; use app::{App, MessageFromUi, MessageToUi, StatefulList, UiFocus};
use config::Config; use config::Config;
use rpc::RpcClient; use rpc::RpcClient;
use tracing::{error, info, warn};
static CONFIG: OnceLock<Config> = OnceLock::new(); static CONFIG: OnceLock<Config> = 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<tracing_appender::non_blocking::WorkerGuard> {
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] #[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> { async fn main() -> Result<(), Box<dyn std::error::Error>> {
let _log_guard = init_tracing();
let config = CONFIG.get_or_init(|| crabidy_core::init_config("cbd-tui.toml")); let config = CONFIG.get_or_init(|| crabidy_core::init_config("cbd-tui.toml"));
let (ui_tx, rx): (Sender<MessageFromUi>, Receiver<MessageFromUi>) = flume::unbounded(); let (ui_tx, rx): (Sender<MessageFromUi>, Receiver<MessageFromUi>) = flume::unbounded();
@ -52,6 +86,7 @@ async fn orchestrate<'a>(
config: &'static Config, config: &'static Config,
(tx, rx): (Sender<MessageToUi>, Receiver<MessageFromUi>), (tx, rx): (Sender<MessageToUi>, Receiver<MessageFromUi>),
) -> Result<(), Box<dyn Error>> { ) -> Result<(), Box<dyn Error>> {
info!(address = config.server.address, "connecting to server");
let mut rpc_client = rpc::RpcClient::connect(&config.server.address).await?; let mut rpc_client = rpc::RpcClient::connect(&config.server.address).await?;
if let Some(root_node) = rpc_client.get_library_node("node:/").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?; let init_data = rpc_client.init().await?;
info!("received initial state from server");
tx.send_async(MessageToUi::Init(init_data)).await?; tx.send_async(MessageToUi::Init(init_data)).await?;
loop { loop {
if let Err(er) = poll(&mut rpc_client, &rx, &tx).await { if let Err(err) = poll(&mut rpc_client, &rx, &tx).await {
println!("ERROR"); error!("request to server failed: {err}");
} }
} }
} }
@ -135,8 +171,10 @@ async fn poll(
tx.send_async(MessageToUi::Update(update)).await?; tx.send_async(MessageToUi::Update(update)).await?;
} }
} }
Err(_) => { Err(err) => {
warn!("update stream broke, reconnecting: {err}");
rpc_client.reconnect_update_stream().await; rpc_client.reconnect_update_stream().await;
info!("update stream reconnected");
} }
} }

View File

@ -39,6 +39,8 @@ impl std::fmt::Display for ProviderError {
} }
} }
impl std::error::Error for ProviderError {}
impl LibraryNode { impl LibraryNode {
pub fn new() -> Self { pub fn new() -> Self {
Self { Self {

View File

@ -16,10 +16,11 @@ pub struct QueueManager {
impl From<QueueManager> for Queue { impl From<QueueManager> for Queue {
fn from(queue_manager: QueueManager) -> Self { fn from(queue_manager: QueueManager) -> Self {
Self { Self {
// A clock step backwards must not panic the playback loop.
timestamp: queue_manager timestamp: queue_manager
.created_at .created_at
.elapsed() .elapsed()
.expect("failed to get elapsed time") .unwrap_or_default()
.as_secs(), .as_secs(),
current_position: queue_manager.current_position() as u32, current_position: queue_manager.current_position() as u32,
tracks: queue_manager.tracks, tracks: queue_manager.tracks,

View File

@ -3,8 +3,8 @@ use crabidy_core::proto::crabidy::{
crabidy_service_server::CrabidyServiceServer, InitResponse, LibraryNode, PlayState, Track, crabidy_service_server::CrabidyServiceServer, InitResponse, LibraryNode, PlayState, Track,
}; };
use crabidy_core::{ProviderClient, ProviderError}; use crabidy_core::{ProviderClient, ProviderError};
use tracing::{debug_span, error, info, instrument, level_filters, warn, Span}; use tracing::{debug, error, info, instrument, warn, Span};
use tracing_subscriber::{filter::Targets, prelude::*}; use tracing_subscriber::{prelude::*, EnvFilter};
mod playback; mod playback;
use playback::Playback; use playback::Playback;
@ -15,34 +15,17 @@ use rpc::RpcService;
use tonic::{transport::Server, Result}; use tonic::{transport::Server, Result};
const LISTEN_ADDR: &str = "0.0.0.0:50051";
#[tokio::main] #[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> { async fn main() -> Result<(), Box<dyn std::error::Error>> {
if let Err(err) = tracing_log::LogTracer::init_with_filter(log::LevelFilter::Debug) { let _log_guard = init_tracing();
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 (update_tx, _) = tokio::sync::broadcast::channel(2048); let (update_tx, _) = tokio::sync::broadcast::channel(2048);
let orchestrator = ProviderOrchestrator::init("") let orchestrator = ProviderOrchestrator::init("").await.map_err(|err| {
.await error!("failed to init provider orchestrator: {err}");
.expect("failed to init orchestrator"); err
})?;
let playback = Playback::new(update_tx.clone(), orchestrator.provider_tx.clone()); let playback = Playback::new(update_tx.clone(), orchestrator.provider_tx.clone());
@ -52,7 +35,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
std::thread::spawn(|| { std::thread::spawn(|| {
poll_play_bus(player_msg, playback_tx); poll_play_bus(player_msg, playback_tx);
}); });
info!("gstreamer bus handler started"); info!("player message forwarder started");
let crabidy_service = RpcService::new( let crabidy_service = RpcService::new(
update_tx, update_tx,
@ -64,7 +47,8 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
playback.run(); playback.run();
info!("playback started"); info!("playback started");
let addr = "0.0.0.0:50051".parse()?; let addr = LISTEN_ADDR.parse()?;
info!(%addr, "grpc server listening");
Server::builder() Server::builder()
.add_service(CrabidyServiceServer::new(crabidy_service)) .add_service(CrabidyServiceServer::new(crabidy_service))
.serve(addr) .serve(addr)
@ -73,165 +57,216 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
Ok(()) 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))] #[instrument(skip(rx, tx))]
fn poll_play_bus(rx: flume::Receiver<PlayerMessage>, tx: flume::Sender<PlaybackMessage>) { fn poll_play_bus(rx: flume::Receiver<PlayerMessage>, tx: flume::Sender<PlaybackMessage>) {
for msg in rx.iter() { for msg in rx.iter() {
let span = debug_span!("play-chan"); let command = match msg {
match msg {
PlayerMessage::EndOfStream => { PlayerMessage::EndOfStream => {
if let Err(err) = tx.send(PlaybackMessage::Next { span }) { debug!("player reported end of stream");
error!("failed to send next message: {}", err); PlaybackCommand::Next
} }
} PlayerMessage::Stopped => PlaybackCommand::StateChanged {
PlayerMessage::Stopped => {
if let Err(err) = tx.send(PlaybackMessage::StateChanged {
state: PlayState::Stopped, state: PlayState::Stopped,
span, },
}) { PlayerMessage::Paused => PlaybackCommand::StateChanged {
error!("failed to send stopped message: {}", err);
}
}
PlayerMessage::Paused => {
if let Err(err) = tx.send(PlaybackMessage::StateChanged {
state: PlayState::Paused, state: PlayState::Paused,
span, },
}) { PlayerMessage::Playing => PlaybackCommand::StateChanged {
error!("failed to send paused message: {}", err);
}
}
PlayerMessage::Playing => {
if let Err(err) = tx.send(PlaybackMessage::StateChanged {
state: PlayState::Playing, state: PlayState::Playing,
span, },
}) { PlayerMessage::Elapsed { duration, elapsed } => PlaybackCommand::PositionChanged {
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, duration: duration.as_millis() as u32,
position: elapsed.as_millis() as u32, position: elapsed.as_millis() as u32,
span, },
}) { PlayerMessage::Duration { duration } => PlaybackCommand::PositionChanged {
error!("failed to send elapsed message: {}", err);
}
}
PlayerMessage::Duration { duration } => {
if let Err(err) = tx.send(PlaybackMessage::PostitionChanged {
duration: duration.as_millis() as u32, duration: duration.as_millis() as u32,
position: 0, position: 0,
span, },
}) { };
error!("failed to send duration message: {}", err); 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)] #[derive(Debug)]
pub enum ProviderMessage { pub enum ProviderCommand {
GetLibraryNode { GetLibraryNode {
uuid: String, uuid: String,
result_tx: flume::Sender<Result<LibraryNode, ProviderError>>, result_tx: flume::Sender<Result<LibraryNode, ProviderError>>,
span: Span,
}, },
GetTrack { GetTrack {
uuid: String, uuid: String,
result_tx: flume::Sender<Result<Track, ProviderError>>, result_tx: flume::Sender<Result<Track, ProviderError>>,
span: Span,
}, },
GetTrackUrls { GetTrackUrls {
uuid: String, uuid: String,
result_tx: flume::Sender<Result<Vec<String>, ProviderError>>, result_tx: flume::Sender<Result<Vec<String>, ProviderError>>,
span: Span,
}, },
FlattenNode { FlattenNode {
uuid: String, uuid: String,
result_tx: flume::Sender<Vec<Track>>, result_tx: flume::Sender<Vec<Track>>,
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)] #[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 { Init {
result_tx: flume::Sender<InitResponse>, result_tx: flume::Sender<InitResponse>,
span: Span,
}, },
Replace { Replace {
uuids: Vec<String>, uuids: Vec<String>,
span: Span,
}, },
Queue { Queue {
uuids: Vec<String>, uuids: Vec<String>,
span: Span,
}, },
Append { Append {
uuids: Vec<String>, uuids: Vec<String>,
span: Span,
}, },
Remove { Remove {
positions: Vec<u32>, positions: Vec<u32>,
span: Span,
}, },
Insert { Insert {
position: u32, position: u32,
uuids: Vec<String>, uuids: Vec<String>,
span: Span,
}, },
Clear { Clear {
exclude_current: bool, exclude_current: bool,
span: Span,
}, },
SetCurrent { SetCurrent {
position: u32, position: u32,
span: Span,
},
ToggleShuffle {
span: Span,
},
ToggleRepeat {
span: Span,
},
TogglePlay {
span: Span,
},
Stop {
span: Span,
}, },
ToggleShuffle,
ToggleRepeat,
TogglePlay,
Stop,
ChangeVolume { ChangeVolume {
delta: f32, delta: f32,
span: Span,
},
ToggleMute {
span: Span,
},
Next {
span: Span,
},
Prev {
span: Span,
},
RestartTrack {
span: Span,
}, },
ToggleMute,
Next,
Prev,
RestartTrack,
StateChanged { StateChanged {
state: PlayState, state: PlayState,
span: Span,
}, },
VolumeChanged { VolumeChanged {
volume: f32, volume: f32,
span: Span,
}, },
MuteChanged { MuteChanged {
muted: bool, muted: bool,
span: Span,
}, },
PostitionChanged { PositionChanged {
duration: u32, duration: u32,
position: 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",
}
}
}

View File

@ -1,5 +1,4 @@
use crate::PlaybackMessage; use crate::{PlaybackCommand, PlaybackMessage, ProviderCommand, ProviderMessage};
use crate::ProviderMessage;
use audio_player::Player; use audio_player::Player;
use crabidy_core::proto::crabidy::QueueModifiers; use crabidy_core::proto::crabidy::QueueModifiers;
use crabidy_core::proto::crabidy::{ use crabidy_core::proto::crabidy::{
@ -9,8 +8,7 @@ use crabidy_core::proto::crabidy::{
use crabidy_core::ProviderError; use crabidy_core::ProviderError;
use crabidy_server::QueueManager; use crabidy_server::QueueManager;
use std::sync::Mutex; use std::sync::Mutex;
use tracing::debug_span; use tracing::{debug, debug_span, error, instrument, trace, warn, Instrument};
use tracing::{debug, error, instrument, trace, warn, Instrument};
pub struct Playback { pub struct Playback {
update_tx: tokio::sync::broadcast::Sender<StreamUpdate>, update_tx: tokio::sync::broadcast::Sender<StreamUpdate>,
@ -27,7 +25,7 @@ impl Playback {
update_tx: tokio::sync::broadcast::Sender<StreamUpdate>, update_tx: tokio::sync::broadcast::Sender<StreamUpdate>,
provider_tx: flume::Sender<ProviderMessage>, provider_tx: flume::Sender<ProviderMessage>,
) -> Self { ) -> Self {
let (playback_tx, playback_rx) = flume::bounded(10); let (playback_tx, playback_rx) = flume::bounded(64);
let queue = Mutex::new(QueueManager::new()); let queue = Mutex::new(QueueManager::new());
let state = Mutex::new(PlayState::Stopped); let state = Mutex::new(PlayState::Stopped);
let player = Player::default(); let player = Player::default();
@ -44,42 +42,41 @@ impl Playback {
pub fn run(self) { pub fn run(self) {
tokio::spawn(async move { tokio::spawn(async move {
while let Ok(message) = self.playback_rx.recv_async().await { while let Ok(PlaybackMessage { span, command }) = self.playback_rx.recv_async().await {
match message { // Attribute all handler events to a span that is a child of
PlaybackMessage::Init { result_tx, span } => { // the span that was current when the command was sent.
let _e = span.enter(); let handler_span =
let repeat; debug_span!(parent: &span, "playback_command", command = command.name());
let shuffle; self.handle_command(command).instrument(handler_span).await;
}
warn!("playback message channel closed, loop exiting");
});
}
async fn handle_command(&self, command: PlaybackCommand) {
trace!("handling playback command");
match command {
PlaybackCommand::Init { result_tx } => {
let response = { let response = {
let Ok(queue) = self.queue.lock() else { let Ok(queue) = self.queue.lock() else {
error!("failed to get queue lock"); error!("queue lock poisoned");
continue; return;
}; };
debug!("got queue lock");
repeat = queue.repeat;
shuffle = queue.shuffle;
let queue_track = QueueTrack { let queue_track = QueueTrack {
queue_position: queue.current_position() as u32, queue_position: queue.current_position() as u32,
track: queue.current_track(), track: queue.current_track(),
}; };
trace!("queue_track {:?}", queue_track);
debug!("released queue_track lock");
let position = TrackPosition { let position = TrackPosition {
duration: 0, duration: 0,
position: 0, position: 0,
}; };
trace!("position {:?}", position);
let play_state = { let play_state = {
debug!("getting play state lock");
let Ok(play_state) = self.state.lock() else { let Ok(play_state) = self.state.lock() else {
error!("failed to get play state lock"); error!("play state lock poisoned");
continue; return;
}; };
*play_state *play_state
}; };
trace!("play_state {:?}", play_state);
debug!("released play state lock");
InitResponse { InitResponse {
queue: Some(queue.clone().into()), queue: Some(queue.clone().into()),
queue_track: Some(queue_track), queue_track: Some(queue_track),
@ -87,590 +84,423 @@ impl Playback {
volume: 0.0, volume: 0.0,
mute: false, mute: false,
position: Some(position), position: Some(position),
mods: Some(QueueModifiers { repeat, shuffle }), mods: Some(QueueModifiers {
repeat: queue.repeat,
shuffle: queue.shuffle,
}),
} }
}; };
trace!("response {:?}", response); trace!(?response, "sending init response");
if let Err(err) = result_tx.send(response) { if let Err(err) = result_tx.send(response) {
error!("failed to send response: {:#?}", err); error!("failed to send init response: {err}");
} }
} }
PlaybackMessage::Replace { uuids, span } => {
let _e = span.enter(); PlaybackCommand::Replace { uuids } => {
let mut all_tracks = Vec::new(); let all_tracks = self.resolve_tracks(uuids).await;
for uuid in uuids { debug!(count = all_tracks.len(), "replacing queue");
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 current = {
let Ok(mut queue) = self.queue.lock() else { let Ok(mut queue) = self.queue.lock() else {
debug!("got queue lock"); error!("queue lock poisoned");
continue; return;
}; };
debug!("got queue lock");
queue.replace_with_tracks(&all_tracks); queue.replace_with_tracks(&all_tracks);
let queue_update_tx = self.update_tx.clone(); self.broadcast(StreamUpdate::Queue(queue.clone().into()));
let update = StreamUpdate::Queue(queue.clone().into());
if let Err(err) = queue_update_tx.send(update) {
trace!("{:?}", err)
};
queue.current_track() queue.current_track()
}; };
debug!("got current {:?}", current); self.play(current).await;
self.play(current).in_current_span().await;
} }
PlaybackMessage::Queue { uuids, span } => { PlaybackCommand::Queue { uuids } => {
let _e = span.enter(); let all_tracks = self.resolve_tracks(uuids).await;
debug!("queing"); debug!(count = all_tracks.len(), "queueing after current");
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 track = {
let Ok(mut queue) = self.queue.lock() else { let Ok(mut queue) = self.queue.lock() else {
debug!("got queue lock"); error!("queue lock poisoned");
continue; return;
}; };
debug!("got queue lock");
let track = queue.queue_tracks(&all_tracks); let track = queue.queue_tracks(&all_tracks);
let queue_update_tx = self.update_tx.clone(); self.broadcast(StreamUpdate::Queue(queue.clone().into()));
let update = StreamUpdate::Queue(queue.clone().into());
if let Err(err) = queue_update_tx.send(update) {
trace!("{:?}", err)
}
track track
}; };
debug!("que lock released"); self.play_if_some(track).await;
self.play(track).in_current_span().await;
} }
PlaybackMessage::Append { uuids, span } => { PlaybackCommand::Append { uuids } => {
let _e = span.enter(); let all_tracks = self.resolve_tracks(uuids).await;
debug!("appending"); debug!(count = all_tracks.len(), "appending to queue");
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 track = {
let Ok(mut queue) = self.queue.lock() else { let Ok(mut queue) = self.queue.lock() else {
debug!("got queue lock"); error!("queue lock poisoned");
continue; return;
}; };
debug!("got queue lock");
let track = queue.append_tracks(&all_tracks); let track = queue.append_tracks(&all_tracks);
let queue_update_tx = self.update_tx.clone(); self.broadcast(StreamUpdate::Queue(queue.clone().into()));
let update = StreamUpdate::Queue(queue.clone().into());
if let Err(err) = queue_update_tx.send(update) {
trace!("{:?}", err)
}
track track
}; };
debug!("queue lock released"); self.play_if_some(track).await;
self.play(track).in_current_span().await;
} }
PlaybackMessage::Remove { positions, span } => { PlaybackCommand::Remove { positions } => {
let _e = span.enter(); debug!(?positions, "removing tracks");
let is_last; let (track, was_last) = {
debug!("removing");
let track = {
let Ok(mut queue) = self.queue.lock() else { let Ok(mut queue) = self.queue.lock() else {
debug!("got queue lock"); error!("queue lock poisoned");
continue; return;
}; };
debug!("got queue lock"); let was_last = queue.is_last_track();
is_last = queue.is_last_track();
let track = queue.remove_tracks(&positions); let track = queue.remove_tracks(&positions);
let queue_update_tx = self.update_tx.clone(); self.broadcast(StreamUpdate::Queue(queue.clone().into()));
let update = StreamUpdate::Queue(queue.clone().into()); (track, was_last)
if let Err(err) = queue_update_tx.send(update) {
trace!("{:?}", err)
}; };
track
};
debug!("queue lock released");
let state = { let state = {
let Ok(state) = self.state.lock() else { let Ok(state) = self.state.lock() else {
error!("failed to get play state lock"); error!("play state lock poisoned");
continue; return;
}; };
*state *state
}; };
if state == PlayState::Playing { if state == PlayState::Playing && track.is_some() {
if is_last { // The playing track was removed: play its successor, or
if let Err(err) = self.player.stop().in_current_span().await { // stop when it was the last one.
error!("{:?}", err) if was_last {
} self.stop_player().await;
} else { } else {
self.play(track).in_current_span().await; self.play(track).await;
} }
} }
} }
PlaybackMessage::Insert { PlaybackCommand::Insert { position, uuids } => {
position, let all_tracks = self.resolve_tracks(uuids).await;
uuids, debug!(count = all_tracks.len(), position, "inserting into queue");
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 track = {
let Ok(mut queue) = self.queue.lock() else { let Ok(mut queue) = self.queue.lock() else {
debug!("got queue lock"); error!("queue lock poisoned");
continue; return;
}; };
debug!("got queue lock");
let track = queue.insert_tracks(position, &all_tracks); let track = queue.insert_tracks(position, &all_tracks);
let queue_update_tx = self.update_tx.clone(); self.broadcast(StreamUpdate::Queue(queue.clone().into()));
let update = StreamUpdate::Queue(queue.clone().into());
if let Err(err) = queue_update_tx.send(update) {
trace!("{:?}", err)
};
track track
}; };
debug!("queue lock released"); self.play_if_some(track).await;
self.play(track).in_current_span().await;
} }
PlaybackMessage::Clear { PlaybackCommand::Clear { exclude_current } => {
exclude_current, debug!(exclude_current, "clearing queue");
span,
} => {
let _e = span.enter();
debug!("clearing");
let should_stop = { let should_stop = {
let Ok(mut queue) = self.queue.lock() else { let Ok(mut queue) = self.queue.lock() else {
debug!("got queue lock"); error!("queue lock poisoned");
continue; return;
}; };
debug!("got queue lock");
let should_stop = queue.clear(exclude_current); let should_stop = queue.clear(exclude_current);
let queue_update_tx = self.update_tx.clone(); self.broadcast(StreamUpdate::Queue(queue.clone().into()));
let update = StreamUpdate::Queue(queue.clone().into());
if let Err(err) = queue_update_tx.send(update) {
trace!("{:?}", err)
};
should_stop should_stop
}; };
debug!("queue lock released");
if should_stop { if should_stop {
if let Err(err) = self.player.stop().in_current_span().await { self.stop_player().await;
error!("{:?}", err)
}
} }
} }
PlaybackMessage::SetCurrent { PlaybackCommand::SetCurrent { position } => {
position: queue_position, debug!(position, "jumping to queue position");
span,
} => {
let _e = span.enter();
debug!("setting current");
let track = { let track = {
let Ok(mut queue) = self.queue.lock() else { let Ok(mut queue) = self.queue.lock() else {
debug!("got queue lock"); error!("queue lock poisoned");
continue; return;
}; };
debug!("got queue lock"); queue.set_current_position(position);
queue.set_current_position(queue_position);
queue.current_track() queue.current_track()
}; };
debug!("quue lock released and got current {:?}", track); self.play(track).await;
self.play(track).in_current_span().await;
} }
PlaybackMessage::ToggleShuffle { span } => { PlaybackCommand::ToggleShuffle => {
let _e = span.enter(); let (shuffle, repeat) = {
debug!("toggling shuffle");
let shuffle;
let repeat;
{
let Ok(mut queue) = self.queue.lock() else { let Ok(mut queue) = self.queue.lock() else {
debug!("got queue lock"); error!("queue lock poisoned");
continue; return;
}; };
debug!("got queue lock");
repeat = queue.repeat;
if queue.shuffle { if queue.shuffle {
queue.shuffle_off() queue.shuffle_off()
} else { } else {
queue.shuffle_on() queue.shuffle_on()
} }
shuffle = queue.shuffle; (queue.shuffle, 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::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"); debug!(shuffle, "toggled shuffle");
shuffle = queue.shuffle; self.broadcast(StreamUpdate::Mods(QueueModifiers { shuffle, repeat }));
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 } => { PlaybackCommand::ToggleRepeat => {
let _e = span.enter(); let (shuffle, repeat) = {
debug!("toggling play"); 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 state = {
let Ok(state) = self.state.lock() else { let Ok(state) = self.state.lock() else {
debug!("got state lock"); error!("play state lock poisoned");
continue; return;
}; };
*state *state
}; };
debug!("got state lock"); debug!(?state, "toggling play");
if state == PlayState::Playing { if state == PlayState::Playing {
if let Err(err) = self.player.pause().await { if let Err(err) = self.player.pause().await {
error!("{:?}", err) warn!("pause failed: {err:?}");
} }
} else if let Err(err) = self.player.unpause().await { } else if let Err(err) = self.player.unpause().await {
error!("{:?}", err) warn!("unpause failed: {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 } => { PlaybackCommand::Stop => {
let _e = span.enter(); debug!("stopping playback");
debug!("changing volume"); self.stop_player().await;
if let Ok(volume) = self.player.volume().await { }
debug!("got volume {:?}", volume);
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 { if let Err(err) = self.player.set_volume(volume + delta).await {
error!("{:?}", err) 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)");
} }
PlaybackMessage::ToggleMute { span } => { PlaybackCommand::Next => {
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 track = {
let Ok(mut queue) = self.queue.lock() else { let Ok(mut queue) = self.queue.lock() else {
debug!("got queue lock"); error!("queue lock poisoned");
continue; return;
}; };
debug!("got queue lock");
queue.next_track() queue.next_track()
}; };
debug!("released queue lock and got track {:?}", track); debug!(
track = track.as_ref().map(|t| t.uuid.as_str()),
self.play_or_stop(track).in_current_span().await; "advancing to next track"
);
self.play_or_stop(track).await;
} }
PlaybackMessage::Prev { span } => { PlaybackCommand::Prev => {
let _e = span.enter();
debug!("preving");
let track = { let track = {
let Ok(mut queue) = self.queue.lock() else { let Ok(mut queue) = self.queue.lock() else {
debug!("got queue lock"); error!("queue lock poisoned");
continue; return;
}; };
debug!("got queue lock");
queue.prev_track() queue.prev_track()
}; };
debug!("released queue lock and got track {:?}", track); debug!(
self.play_or_stop(track).in_current_span().await; track = track.as_ref().map(|t| t.uuid.as_str()),
"going back to previous track"
);
self.play_or_stop(track).await;
} }
PlaybackMessage::StateChanged { state, span } => { PlaybackCommand::StateChanged { state } => {
let _e = span.enter();
debug!("state changed");
let play_state = { let play_state = {
let Ok(mut state_lock) = self.state.lock() else { let Ok(mut state_lock) = self.state.lock() else {
debug!("got state lock"); error!("play state lock poisoned");
continue; return;
}; };
*state_lock = state; *state_lock = state;
state state
}; };
debug!("released state lock and got play state {:?}", play_state); debug!(?play_state, "player state changed");
let active_track_tx = self.update_tx.clone(); self.broadcast(StreamUpdate::PlayState(play_state as i32));
let update = StreamUpdate::PlayState(play_state as i32);
if let Err(err) = active_track_tx.send(update) {
trace!("{:?}", err)
};
} }
PlaybackMessage::RestartTrack { span } => { PlaybackCommand::RestartTrack => {
let _e = span.enter(); debug!("restarting current track");
debug!("restarting track");
if let Err(err) = self.player.restart().await { if let Err(err) = self.player.restart().await {
error!("{:?}", err) warn!("restart failed: {err:?}");
} }
} }
PlaybackMessage::VolumeChanged { volume, span } => { PlaybackCommand::VolumeChanged { volume } => {
let _e = span.enter(); trace!(volume, "volume changed");
trace!("volume changed"); self.broadcast(StreamUpdate::Volume(volume));
let update_tx = self.update_tx.clone(); }
let update = StreamUpdate::Volume(volume);
if let Err(err) = update_tx.send(update) { PlaybackCommand::MuteChanged { muted } => {
trace!("{:?}", err) 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 }));
}
} }
} }
PlaybackMessage::MuteChanged { muted, span } => { /// Sends an update to all connected clients. Having no subscribers is
let _e = span.enter(); /// normal and not an error.
trace!("mute changed"); fn broadcast(&self, update: StreamUpdate) {
let update_tx = self.update_tx.clone(); if let Err(err) = self.update_tx.send(update) {
let update = StreamUpdate::Mute(muted); trace!("no update stream subscribers: {err}");
if let Err(err) = update_tx.send(update) {
trace!("{:?}", err)
} }
} }
PlaybackMessage::PostitionChanged { /// Resolves a mixed list of track and node identifiers into tracks by
duration, /// asking the provider orchestrator.
position, async fn resolve_tracks(&self, uuids: Vec<String>) -> Vec<Track> {
span, let mut all_tracks = Vec::new();
} => { for uuid in uuids {
let _e = span.enter(); if is_track(&uuid) {
trace!("position changed"); match self.get_track(&uuid).await {
let update_tx = self.update_tx.clone(); Ok(track) => all_tracks.push(track),
let update = StreamUpdate::Position(TrackPosition { duration, position }); Err(err) => warn!(uuid, "failed to resolve track: {err}"),
if let Err(err) = update_tx.send(update) { }
trace!("{:?}", 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))] #[instrument(skip(self))]
async fn flatten_node(&self, uuid: &str) -> Vec<Track> { async fn flatten_node(&self, uuid: &str) -> Vec<Track> {
debug!("flattening node");
let tx = self.provider_tx.clone();
let (result_tx, result_rx) = flume::bounded(1); let (result_tx, result_rx) = flume::bounded(1);
let span = debug_span!("prov-chan"); let message = ProviderMessage::new(ProviderCommand::FlattenNode {
let Ok(_) = tx
.send_async(ProviderMessage::FlattenNode {
uuid: uuid.to_string(), uuid: uuid.to_string(),
result_tx, result_tx,
span, });
}) if let Err(err) = self.provider_tx.send_async(message).await {
.in_current_span() error!("provider channel closed: {err}");
.await
else {
return Vec::new(); return Vec::new();
}; }
let Ok(tracks) = result_rx.recv_async().in_current_span().await else { match result_rx.recv_async().await {
return Vec::new(); Ok(tracks) => tracks,
}; Err(err) => {
tracks error!("provider dropped flatten_node reply: {err}");
Vec::new()
}
}
} }
#[instrument(skip(self))] #[instrument(skip(self))]
async fn get_track(&self, uuid: &str) -> Result<Track, ProviderError> { async fn get_track(&self, uuid: &str) -> Result<Track, ProviderError> {
debug!("getting track");
let tx = self.provider_tx.clone();
let (result_tx, result_rx) = flume::bounded(1); let (result_tx, result_rx) = flume::bounded(1);
let span = tracing::trace_span!("prov-chan"); let message = ProviderMessage::new(ProviderCommand::GetTrack {
tx.send_async(ProviderMessage::GetTrack {
uuid: uuid.to_string(), uuid: uuid.to_string(),
result_tx, result_tx,
span, });
}) self.provider_tx
.in_current_span() .send_async(message)
.await .await
.map_err(|_| ProviderError::InternalError)?; .map_err(|_| ProviderError::InternalError)?;
result_rx result_rx
.recv_async() .recv_async()
.in_current_span()
.await .await
.map_err(|_| ProviderError::InternalError)? .map_err(|_| ProviderError::InternalError)?
} }
#[instrument(skip(self))] #[instrument(skip(self))]
async fn get_urls_for_track(&self, uuid: &str) -> Result<Vec<String>, ProviderError> { async fn get_urls_for_track(&self, uuid: &str) -> Result<Vec<String>, ProviderError> {
debug!("getting urls for track");
let tx = self.provider_tx.clone();
let (result_tx, result_rx) = flume::bounded(1); let (result_tx, result_rx) = flume::bounded(1);
let span = tracing::trace_span!("prov-chan"); let message = ProviderMessage::new(ProviderCommand::GetTrackUrls {
tx.send_async(ProviderMessage::GetTrackUrls {
uuid: uuid.to_string(), uuid: uuid.to_string(),
result_tx, result_tx,
span, });
}) self.provider_tx
.in_current_span() .send_async(message)
.await .await
.map_err(|_| ProviderError::InternalError)?; .map_err(|_| ProviderError::InternalError)?;
result_rx result_rx
.recv_async() .recv_async()
.in_current_span()
.await .await
.map_err(|_| ProviderError::InternalError)? .map_err(|_| ProviderError::InternalError)?
} }
#[instrument(skip(self))] async fn stop_player(&self) {
async fn play_or_stop(&self, track: Option<Track>) { if let Err(err) = self.player.stop().await {
debug!("play or stop"); debug!("stop had no effect: {err:?}");
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)
} }
} }
#[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<Track>) {
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<Track>) {
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<Track>) { async fn play(&self, track: Option<Track>) {
debug!("play"); let Some(track) = track else {
if let Some(track) = track { debug!("nothing to play");
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; return;
}; };
if let Some(track) = queue.next_track() { let mut uuid = track.uuid.clone();
track.uuid.clone() let urls = loop {
} else { 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;
};
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; return;
}
}
} }
} }
}; };
{ {
let Ok(queue) = self.queue.lock() else { let Ok(queue) = self.queue.lock() else {
error!("poisend queue lock"); error!("queue lock poisoned");
return; return;
}; };
let queue_update_tx = self.update_tx.clone(); self.broadcast(StreamUpdate::QueueTrack(QueueTrack {
let track = queue.current_track();
let update = StreamUpdate::QueueTrack(QueueTrack {
queue_position: queue.current_position() as u32, queue_position: queue.current_position() as u32,
track, track: queue.current_track(),
}); }));
if let Err(err) = queue_update_tx.send(update) {
trace!("{:?}", err)
}
} }
debug!(url_count = urls.len(), "starting player");
if let Err(err) = self.player.play(&urls[0]).await { if let Err(err) = self.player.play(&urls[0]).await {
error!("{:?}", err) error!("player failed to start track: {err:?}");
}
} }
} }
} }

View File

@ -1,88 +1,80 @@
use crate::ProviderMessage; use crate::{ProviderCommand, ProviderMessage};
use async_trait::async_trait; use async_trait::async_trait;
use crabidy_core::{ use crabidy_core::{
proto::crabidy::{LibraryNode, LibraryNodeChild, Track}, proto::crabidy::{LibraryNode, LibraryNodeChild, Track},
ProviderClient, ProviderError, ProviderClient, ProviderError,
}; };
use std::{fs, path::PathBuf, sync::Arc}; 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)] #[derive(Debug)]
pub struct ProviderOrchestrator { pub struct ProviderOrchestrator {
pub provider_tx: flume::Sender<ProviderMessage>, pub provider_tx: flume::Sender<ProviderMessage>,
provider_rx: flume::Receiver<ProviderMessage>, provider_rx: flume::Receiver<ProviderMessage>,
// known_tracks: RwLock<HashMap<String, Track>>,
// known_nodes: RwLock<HashMap<String, LibraryNode>>,
tidal_client: Arc<tidaldy::Client>, tidal_client: Arc<tidaldy::Client>,
} }
impl ProviderOrchestrator { impl ProviderOrchestrator {
pub fn run(self) { pub fn run(self) {
tokio::spawn(async move { tokio::spawn(async move {
while let Ok(msg) = self.provider_rx.recv_async().await { while let Ok(ProviderMessage { span, command }) = self.provider_rx.recv_async().await {
match msg { let handler_span =
ProviderMessage::GetLibraryNode { debug_span!(parent: &span, "provider_command", command = command.name());
uuid, self.handle_command(command).instrument(handler_span).await;
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);
}
}
}
} }
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))] #[instrument(skip(self))]
async fn flatten_node(&self, node_uuid: &str) -> Vec<Track> { async fn flatten_node(&self, node_uuid: &str) -> Vec<Track> {
let mut tracks = Vec::with_capacity(1000); let mut tracks = Vec::new();
let mut nodes_to_go = Vec::with_capacity(100); let mut nodes_to_go = vec![node_uuid.to_string()];
nodes_to_go.push(node_uuid.to_string());
while let Some(node_uuid) = nodes_to_go.pop() { while let Some(node_uuid) = nodes_to_go.pop() {
let Ok(node) = self.get_lib_node(&node_uuid).in_current_span().await else { let node = match self.get_lib_node(&node_uuid).await {
Ok(node) => node,
Err(err) => {
warn!(node = node_uuid, "skipping unreadable node: {err}");
continue; continue;
}
}; };
if node.is_queable { if node.is_queable {
tracks.extend(node.tracks); tracks.extend(node.tracks);
nodes_to_go.extend(node.children.into_iter().map(|c| c.uuid)) nodes_to_go.extend(node.children.into_iter().map(|c| c.uuid))
} }
} }
debug!(count = tracks.len(), "flattened node into tracks");
tracks tracks
} }
} }
@ -95,29 +87,25 @@ impl ProviderClient for ProviderOrchestrator {
.map(|d| d.join("crabidy")) .map(|d| d.join("crabidy"))
.unwrap_or(PathBuf::from("/tmp")); .unwrap_or(PathBuf::from("/tmp"));
let dir_exists = tokio::fs::try_exists(&config_dir) let dir_exists = tokio::fs::try_exists(&config_dir)
.in_current_span()
.await .await
.map_err(|e| ProviderError::Config(e.to_string()))?; .map_err(|e| ProviderError::Config(e.to_string()))?;
if !dir_exists { if !dir_exists {
tokio::fs::create_dir(&config_dir) tokio::fs::create_dir(&config_dir)
.in_current_span()
.await .await
.map_err(|e| ProviderError::Config(e.to_string()))?; .map_err(|e| ProviderError::Config(e.to_string()))?;
} }
let config_file = config_dir.join("tidaly.toml"); let config_file = config_dir.join("tidaly.toml");
let raw_toml_settings = fs::read_to_string(&config_file).unwrap_or("".to_owned()); debug!(config_file = %config_file.display(), "loading tidal config");
let tidal_client = Arc::new( let raw_toml_settings = fs::read_to_string(&config_file).unwrap_or_default();
tidaldy::Client::init(&raw_toml_settings) let tidal_client = Arc::new(tidaldy::Client::init(&raw_toml_settings).await.map_err(
.in_current_span() |err| {
.await error!("failed to init tidal client: {err}");
.expect("Failed to init Tidal clienta"), err
); },
)?);
let new_toml_config = tidal_client.settings(); let new_toml_config = tidal_client.settings();
if let Err(err) = tokio::fs::write(&config_file, new_toml_config) if let Err(err) = tokio::fs::write(&config_file, new_toml_config).await {
.in_current_span() error!("failed to write tidal config file: {err}");
.await
{
error!("Failed to write config file: {}", err);
}; };
let (provider_tx, provider_rx) = flume::bounded(100); let (provider_tx, provider_rx) = flume::bounded(100);
Ok(Self { Ok(Self {
@ -126,46 +114,38 @@ impl ProviderClient for ProviderOrchestrator {
tidal_client, tidal_client,
}) })
} }
#[instrument(skip(self))]
fn settings(&self) -> String { fn settings(&self) -> String {
"".to_owned() String::new()
} }
#[instrument(skip(self))] #[instrument(skip(self))]
async fn get_urls_for_track(&self, track_uuid: &str) -> Result<Vec<String>, ProviderError> { async fn get_urls_for_track(&self, track_uuid: &str) -> Result<Vec<String>, ProviderError> {
debug!("get_urls_for_track"); self.tidal_client.get_urls_for_track(track_uuid).await
self.tidal_client
.get_urls_for_track(track_uuid)
.in_current_span()
.await
} }
#[instrument(skip(self))] #[instrument(skip(self))]
async fn get_metadata_for_track(&self, track_uuid: &str) -> Result<Track, ProviderError> { async fn get_metadata_for_track(&self, track_uuid: &str) -> Result<Track, ProviderError> {
debug!("get_metadata_for_track"); self.tidal_client.get_metadata_for_track(track_uuid).await
self.tidal_client
.get_metadata_for_track(track_uuid)
.in_current_span()
.await
} }
#[instrument(skip(self))]
fn get_lib_root(&self) -> LibraryNode { fn get_lib_root(&self) -> LibraryNode {
debug!("get_lib_root in provider manager");
let mut root_node = LibraryNode::new(); let mut root_node = LibraryNode::new();
let child = LibraryNodeChild::new("node:tidal".to_owned(), "tidal".to_owned(), false); let child = LibraryNodeChild::new("node:tidal".to_owned(), "tidal".to_owned(), false);
root_node.children.push(child); root_node.children.push(child);
root_node root_node
} }
#[instrument(skip(self))] #[instrument(skip(self))]
async fn get_lib_node(&self, uuid: &str) -> Result<LibraryNode, ProviderError> { async fn get_lib_node(&self, uuid: &str) -> Result<LibraryNode, ProviderError> {
debug!("get_lib_node in provider manager");
if uuid == "node:/" { if uuid == "node:/" {
debug!("get global root"); debug!("serving global library root");
return Ok(self.get_lib_root()); return Ok(self.get_lib_root());
} }
if uuid == "node:tidal" { if uuid == "node:tidal" {
debug!("get tidal root"); debug!("serving tidal library root");
return Ok(self.tidal_client.get_lib_root()); return Ok(self.tidal_client.get_lib_root());
} }
debug!("tidal node"); self.tidal_client.get_lib_node(uuid).await
self.tidal_client.get_lib_node(uuid).in_current_span().await
} }
} }

View File

@ -1,4 +1,4 @@
use crate::{PlaybackMessage, ProviderMessage}; use crate::{PlaybackCommand, PlaybackMessage, ProviderCommand, ProviderMessage};
use crabidy_core::proto::crabidy::{ use crabidy_core::proto::crabidy::{
crabidy_service_server::CrabidyService, get_update_stream_response::Update as StreamUpdate, crabidy_service_server::CrabidyService, get_update_stream_response::Update as StreamUpdate,
AppendRequest, AppendResponse, ChangeVolumeRequest, ChangeVolumeResponse, ClearQueueRequest, AppendRequest, AppendResponse, ChangeVolumeRequest, ChangeVolumeResponse, ClearQueueRequest,
@ -10,11 +10,10 @@ use crabidy_core::proto::crabidy::{
StopResponse, ToggleMuteRequest, ToggleMuteResponse, TogglePlayRequest, TogglePlayResponse, StopResponse, ToggleMuteRequest, ToggleMuteResponse, TogglePlayRequest, TogglePlayResponse,
ToggleRepeatRequest, ToggleRepeatResponse, ToggleShuffleRequest, ToggleShuffleResponse, ToggleRepeatRequest, ToggleRepeatResponse, ToggleShuffleRequest, ToggleShuffleResponse,
}; };
use futures::TryStreamExt;
use std::pin::Pin; use std::pin::Pin;
use tokio_stream::StreamExt; use tokio_stream::StreamExt;
use tonic::{Request, Response, Status}; use tonic::{Request, Response, Status};
use tracing::{debug, debug_span, error, instrument, trace, Instrument, Span}; use tracing::{debug, error, instrument, trace};
#[derive(Debug)] #[derive(Debug)]
pub struct RpcService { pub struct RpcService {
@ -25,16 +24,29 @@ pub struct RpcService {
impl RpcService { impl RpcService {
pub fn new( pub fn new(
update_rx: tokio::sync::broadcast::Sender<StreamUpdate>, update_tx: tokio::sync::broadcast::Sender<StreamUpdate>,
playback_tx: flume::Sender<PlaybackMessage>, playback_tx: flume::Sender<PlaybackMessage>,
provider_tx: flume::Sender<ProviderMessage>, provider_tx: flume::Sender<ProviderMessage>,
) -> Self { ) -> Self {
Self { Self {
update_tx: update_rx, update_tx,
playback_tx, playback_tx,
provider_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] #[tonic::async_trait]
@ -44,25 +56,13 @@ impl CrabidyService for RpcService {
#[instrument(skip(self, _request))] #[instrument(skip(self, _request))]
async fn init(&self, _request: Request<InitRequest>) -> Result<Response<InitResponse>, Status> { async fn init(&self, _request: Request<InitRequest>) -> Result<Response<InitResponse>, Status> {
debug!("Received init request"); debug!("received init request");
let playback_tx = self.playback_tx.clone();
let (result_tx, result_rx) = flume::bounded(1); let (result_tx, result_rx) = flume::bounded(1);
let span = debug_span!("play-chan"); self.send_playback(PlaybackCommand::Init { result_tx })
if let Err(err) = playback_tx .await?;
.send_async(PlaybackMessage::Init { result_tx, span }) let response = result_rx.recv_async().await.map_err(|err| {
.in_current_span() error!("no reply from playback loop: {err}");
.await Status::internal("playback loop did not reply")
{
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")
})?; })?;
Ok(Response::new(response)) Ok(Response::new(response))
} }
@ -73,32 +73,27 @@ impl CrabidyService for RpcService {
request: Request<GetLibraryNodeRequest>, request: Request<GetLibraryNodeRequest>,
) -> Result<Response<GetLibraryNodeResponse>, Status> { ) -> Result<Response<GetLibraryNodeResponse>, Status> {
let uuid = request.into_inner().uuid; let uuid = request.into_inner().uuid;
Span::current().record("uuid", &uuid); tracing::Span::current().record("uuid", uuid.as_str());
debug!("Received get_library_node request"); debug!("received get_library_node request");
let provider_tx = self.provider_tx.clone();
let (result_tx, result_rx) = flume::bounded(1); let (result_tx, result_rx) = flume::bounded(1);
let span = debug_span!("prov-chan"); self.provider_tx
provider_tx .send_async(ProviderMessage::new(ProviderCommand::GetLibraryNode {
.send_async(ProviderMessage::GetLibraryNode {
uuid, uuid,
result_tx, result_tx,
span, }))
})
.in_current_span()
.await .await
.map_err(|_| Status::internal("Failed to send request via channel"))?; .map_err(|err| {
let result = result_rx error!("provider channel closed: {err}");
.recv_async() Status::internal("provider unavailable")
.in_current_span() })?;
.await let result = result_rx.recv_async().await.map_err(|err| {
.map_err(|e| { error!("no reply from provider: {err}");
error!("{:?}", e); Status::internal("provider did not reply")
Status::internal("Failed to receive response from provider channel")
})?; })?;
match result { match result {
Ok(node) => Ok(Response::new(GetLibraryNodeResponse { node: Some(node) })), Ok(node) => Ok(Response::new(GetLibraryNodeResponse { node: Some(node) })),
Err(err) => { Err(err) => {
error!("{:?}", err); error!("get_library_node failed: {err}");
Err(Status::internal(err.to_string())) Err(Status::internal(err.to_string()))
} }
} }
@ -107,204 +102,137 @@ impl CrabidyService for RpcService {
#[instrument(skip(self, request), fields(uuids))] #[instrument(skip(self, request), fields(uuids))]
async fn queue( async fn queue(
&self, &self,
request: tonic::Request<QueueRequest>, request: Request<QueueRequest>,
) -> std::result::Result<tonic::Response<QueueResponse>, tonic::Status> { ) -> Result<Response<QueueResponse>, Status> {
let uuids = request.into_inner().uuids.clone(); let uuids = request.into_inner().uuids;
Span::current().record("uuids", format!("{:?}", uuids)); tracing::Span::current().record("uuids", format!("{uuids:?}"));
debug!("Received queue request"); debug!("received queue request");
let playback_tx = self.playback_tx.clone(); self.send_playback(PlaybackCommand::Queue { uuids }).await?;
let span = debug_span!("play-chan"); Ok(Response::new(QueueResponse {}))
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))
} }
#[instrument(skip(self, request), fields(uuids))] #[instrument(skip(self, request), fields(uuids))]
async fn replace( async fn replace(
&self, &self,
request: tonic::Request<ReplaceRequest>, request: Request<ReplaceRequest>,
) -> std::result::Result<tonic::Response<ReplaceResponse>, tonic::Status> { ) -> Result<Response<ReplaceResponse>, Status> {
let uuids = request.into_inner().uuids.clone(); let uuids = request.into_inner().uuids;
Span::current().record("uuids", format!("{:?}", uuids)); tracing::Span::current().record("uuids", format!("{uuids:?}"));
debug!("Received replace request"); debug!("received replace request");
let playback_tx = self.playback_tx.clone(); self.send_playback(PlaybackCommand::Replace { uuids })
let span = debug_span!("play-chan"); .await?;
playback_tx Ok(Response::new(ReplaceResponse {}))
.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))
} }
#[instrument(skip(self, request), fields(uuids))] #[instrument(skip(self, request), fields(uuids))]
async fn append( async fn append(
&self, &self,
request: tonic::Request<AppendRequest>, request: Request<AppendRequest>,
) -> std::result::Result<tonic::Response<AppendResponse>, tonic::Status> { ) -> Result<Response<AppendResponse>, Status> {
let uuids = request.into_inner().uuids.clone(); let uuids = request.into_inner().uuids;
Span::current().record("uuids", format!("{:?}", uuids)); tracing::Span::current().record("uuids", format!("{uuids:?}"));
debug!("Received append request"); debug!("received append request");
let playback_tx = self.playback_tx.clone(); self.send_playback(PlaybackCommand::Append { uuids })
let span = debug_span!("play-chan"); .await?;
playback_tx Ok(Response::new(AppendResponse {}))
.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))
} }
#[instrument(skip(self, request), fields(positions))] #[instrument(skip(self, request), fields(positions))]
async fn remove( async fn remove(
&self, &self,
request: tonic::Request<RemoveRequest>, request: Request<RemoveRequest>,
) -> std::result::Result<tonic::Response<RemoveResponse>, tonic::Status> { ) -> Result<Response<RemoveResponse>, Status> {
let positions = request.into_inner().positions; let positions = request.into_inner().positions;
Span::current().record("positions", format!("{:?}", positions)); tracing::Span::current().record("positions", format!("{positions:?}"));
debug!("Received remove request"); debug!("received remove request");
let playback_tx = self.playback_tx.clone(); self.send_playback(PlaybackCommand::Remove { positions })
let span = debug_span!("play-chan"); .await?;
playback_tx Ok(Response::new(RemoveResponse {}))
.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))
} }
#[instrument(skip(self, request), fields(uuids, position))] #[instrument(skip(self, request), fields(uuids, position))]
async fn insert( async fn insert(
&self, &self,
request: tonic::Request<InsertRequest>, request: Request<InsertRequest>,
) -> std::result::Result<tonic::Response<InsertResponse>, tonic::Status> { ) -> Result<Response<InsertResponse>, Status> {
let req = request.into_inner(); let req = request.into_inner();
let uuids = req.uuids.clone(); tracing::Span::current().record("uuids", format!("{:?}", req.uuids));
let position = req.position; tracing::Span::current().record("position", req.position);
Span::current().record("uuids", format!("{:?}", uuids)); debug!("received insert request");
Span::current().record("position", position); self.send_playback(PlaybackCommand::Insert {
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, position: req.position,
uuids, uuids: req.uuids,
span,
}) })
.in_current_span() .await?;
.await Ok(Response::new(InsertResponse {}))
.map_err(|_| Status::internal("Failed to send request via channel"))?;
let reply = InsertResponse {};
Ok(Response::new(reply))
} }
#[instrument(skip(self, request), fields(exclude_current))] #[instrument(skip(self, request), fields(exclude_current))]
async fn clear_queue( async fn clear_queue(
&self, &self,
request: tonic::Request<ClearQueueRequest>, request: Request<ClearQueueRequest>,
) -> std::result::Result<tonic::Response<ClearQueueResponse>, tonic::Status> { ) -> Result<Response<ClearQueueResponse>, Status> {
let exclude_current = request.into_inner().exclude_current; let exclude_current = request.into_inner().exclude_current;
Span::current().record("exclude_current", exclude_current); tracing::Span::current().record("exclude_current", exclude_current);
debug!("Received clear_queue request"); debug!("received clear_queue request");
let playback_tx = self.playback_tx.clone(); self.send_playback(PlaybackCommand::Clear { exclude_current })
let span = debug_span!("play-chan"); .await?;
playback_tx Ok(Response::new(ClearQueueResponse {}))
.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))
} }
#[instrument(skip(self, request), fields(position))] #[instrument(skip(self, request), fields(position))]
async fn set_current( async fn set_current(
&self, &self,
request: tonic::Request<SetCurrentRequest>, request: Request<SetCurrentRequest>,
) -> std::result::Result<tonic::Response<SetCurrentResponse>, tonic::Status> { ) -> Result<Response<SetCurrentResponse>, Status> {
let position = request.into_inner().position; let position = request.into_inner().position;
Span::current().record("position", position); tracing::Span::current().record("position", position);
debug!("Received set_current request"); debug!("received set_current request");
let playback_tx = self.playback_tx.clone(); self.send_playback(PlaybackCommand::SetCurrent { position })
let span = debug_span!("play-chan"); .await?;
playback_tx Ok(Response::new(SetCurrentResponse {}))
.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))
} }
#[instrument(skip(self, _request))] #[instrument(skip(self, _request))]
async fn toggle_shuffle( async fn toggle_shuffle(
&self, &self,
_request: tonic::Request<ToggleShuffleRequest>, _request: Request<ToggleShuffleRequest>,
) -> std::result::Result<tonic::Response<ToggleShuffleResponse>, tonic::Status> { ) -> Result<Response<ToggleShuffleResponse>, Status> {
debug!("Received toggle_shuffle request"); debug!("received toggle_shuffle request");
let playback_tx = self.playback_tx.clone(); self.send_playback(PlaybackCommand::ToggleShuffle).await?;
let span = debug_span!("play-chan"); Ok(Response::new(ToggleShuffleResponse {}))
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))
} }
#[instrument(skip(self, _request))] #[instrument(skip(self, _request))]
async fn toggle_repeat( async fn toggle_repeat(
&self, &self,
_request: tonic::Request<ToggleRepeatRequest>, _request: Request<ToggleRepeatRequest>,
) -> std::result::Result<tonic::Response<ToggleRepeatResponse>, tonic::Status> { ) -> Result<Response<ToggleRepeatResponse>, Status> {
debug!("Received toggle_repeat request"); debug!("received toggle_repeat request");
let playback_tx = self.playback_tx.clone(); self.send_playback(PlaybackCommand::ToggleRepeat).await?;
let span = debug_span!("play-chan"); Ok(Response::new(ToggleRepeatResponse {}))
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))
} }
#[instrument(skip(self, _request))] #[instrument(skip(self, _request))]
async fn get_update_stream( async fn get_update_stream(
&self, &self,
_request: tonic::Request<GetUpdateStreamRequest>, _request: Request<GetUpdateStreamRequest>,
) -> std::result::Result<tonic::Response<Self::GetUpdateStreamStream>, tonic::Status> { ) -> Result<Response<Self::GetUpdateStreamStream>, Status> {
debug!("Received get_update_stream request"); debug!("received get_update_stream request, subscribing client");
let update_rx = self.update_tx.subscribe(); let update_rx = self.update_tx.subscribe();
let update_stream = tokio_stream::wrappers::BroadcastStream::new(update_rx); let update_stream = tokio_stream::wrappers::BroadcastStream::new(update_rx);
let output_stream = update_stream.into_stream().map(|update_result| { let output_stream = update_stream.map(|update_result| {
trace!("Got update: {:?}", update_result); trace!(?update_result, "forwarding update");
match update_result { match update_result {
Ok(update) => Ok(GetUpdateStreamResponse { Ok(update) => Ok(GetUpdateStreamResponse {
update: Some(update), update: Some(update),
}), }),
Err(_) => Err(tonic::Status::new( Err(err) => {
tonic::Code::Unknown, // The client lagged too far behind the broadcast channel.
"Internal channel error", error!("update stream lagged: {err}");
)), Err(Status::data_loss("update stream lagged"))
}
} }
}); });
@ -314,146 +242,73 @@ impl CrabidyService for RpcService {
#[instrument(skip(self, _request))] #[instrument(skip(self, _request))]
async fn save_queue( async fn save_queue(
&self, &self,
_request: tonic::Request<SaveQueueRequest>, _request: Request<SaveQueueRequest>,
) -> std::result::Result<tonic::Response<SaveQueueResponse>, tonic::Status> { ) -> Result<Response<SaveQueueResponse>, Status> {
debug!("Received save_queue request"); debug!("received save_queue request (not implemented)");
let reply = SaveQueueResponse {}; Ok(Response::new(SaveQueueResponse {}))
Ok(Response::new(reply))
} }
/// Playback
#[instrument(skip(self, _request))] #[instrument(skip(self, _request))]
async fn toggle_play( async fn toggle_play(
&self, &self,
_request: tonic::Request<TogglePlayRequest>, _request: Request<TogglePlayRequest>,
) -> std::result::Result<tonic::Response<TogglePlayResponse>, tonic::Status> { ) -> Result<Response<TogglePlayResponse>, Status> {
debug!("Received toggle_play request"); debug!("received toggle_play request");
let playback_tx = self.playback_tx.clone(); self.send_playback(PlaybackCommand::TogglePlay).await?;
let span = debug_span!("play-chan"); Ok(Response::new(TogglePlayResponse {}))
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))
} }
#[instrument(skip(self, _request))] #[instrument(skip(self, _request))]
async fn stop( async fn stop(&self, _request: Request<StopRequest>) -> Result<Response<StopResponse>, Status> {
&self, debug!("received stop request");
_request: tonic::Request<StopRequest>, self.send_playback(PlaybackCommand::Stop).await?;
) -> std::result::Result<tonic::Response<StopResponse>, tonic::Status> { Ok(Response::new(StopResponse {}))
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))
} }
#[instrument(skip(self, request), fields(delta))] #[instrument(skip(self, request), fields(delta))]
async fn change_volume( async fn change_volume(
&self, &self,
request: tonic::Request<ChangeVolumeRequest>, request: Request<ChangeVolumeRequest>,
) -> std::result::Result<tonic::Response<ChangeVolumeResponse>, tonic::Status> { ) -> Result<Response<ChangeVolumeResponse>, Status> {
let delta = request.into_inner().delta; let delta = request.into_inner().delta;
Span::current().record("delta", delta); tracing::Span::current().record("delta", delta);
debug!("Received change_volume request"); debug!("received change_volume request");
let playback_tx = self.playback_tx.clone(); self.send_playback(PlaybackCommand::ChangeVolume { delta })
let span = debug_span!("play-chan"); .await?;
if let Err(err) = playback_tx Ok(Response::new(ChangeVolumeResponse {}))
.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))
} }
#[instrument(skip(self, _request))] #[instrument(skip(self, _request))]
async fn toggle_mute( async fn toggle_mute(
&self, &self,
_request: tonic::Request<ToggleMuteRequest>, _request: Request<ToggleMuteRequest>,
) -> std::result::Result<tonic::Response<ToggleMuteResponse>, tonic::Status> { ) -> Result<Response<ToggleMuteResponse>, Status> {
debug!("Received toggle_mute request"); debug!("received toggle_mute request");
let playback_tx = self.playback_tx.clone(); self.send_playback(PlaybackCommand::ToggleMute).await?;
let span = debug_span!("play-chan"); Ok(Response::new(ToggleMuteResponse {}))
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))
} }
#[instrument(skip(self, _request))] #[instrument(skip(self, _request))]
async fn next( async fn next(&self, _request: Request<NextRequest>) -> Result<Response<NextResponse>, Status> {
&self, debug!("received next request");
_request: tonic::Request<NextRequest>, self.send_playback(PlaybackCommand::Next).await?;
) -> std::result::Result<tonic::Response<NextResponse>, tonic::Status> { Ok(Response::new(NextResponse {}))
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))
} }
#[instrument(skip(self, _request))] #[instrument(skip(self, _request))]
async fn prev( async fn prev(&self, _request: Request<PrevRequest>) -> Result<Response<PrevResponse>, Status> {
&self, debug!("received prev request");
_request: tonic::Request<PrevRequest>, self.send_playback(PlaybackCommand::Prev).await?;
) -> std::result::Result<tonic::Response<PrevResponse>, tonic::Status> { Ok(Response::new(PrevResponse {}))
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))
} }
#[instrument(skip(self, _request))] #[instrument(skip(self, _request))]
async fn restart_track( async fn restart_track(
&self, &self,
_request: tonic::Request<RestartTrackRequest>, _request: Request<RestartTrackRequest>,
) -> std::result::Result<tonic::Response<RestartTrackResponse>, tonic::Status> { ) -> Result<Response<RestartTrackResponse>, Status> {
debug!("Received restart_track request"); debug!("received restart_track request");
let playback_tx = self.playback_tx.clone(); self.send_playback(PlaybackCommand::RestartTrack).await?;
let span = debug_span!("play-chan"); Ok(Response::new(RestartTrackResponse {}))
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))
} }
} }

View File

@ -3,7 +3,7 @@
use reqwest::Client as HttpClient; use reqwest::Client as HttpClient;
use serde::de::DeserializeOwned; use serde::de::DeserializeOwned;
use tokio::time::{sleep, Duration, Instant}; 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 config;
pub mod models; pub mod models;
use async_trait::async_trait; 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) { let settings: config::Settings = if let Ok(settings) = toml::from_str(raw_toml_settings) {
settings settings
} else { } else {
let settings = config::Settings::default(); warn!("could not parse toml settings, using defaults");
println!( config::Settings::default()
"could not parse toml settings: {:#?} using default settings instead: {:#?}",
raw_toml_settings, settings
);
settings
}; };
let mut client = Self::new(settings)?; let mut client = Self::new(settings)?;
@ -48,16 +44,21 @@ impl crabidy_core::ProviderClient for Client {
&self, &self,
track_uuid: &str, track_uuid: &str,
) -> Result<Vec<String>, crabidy_core::ProviderError> { ) -> Result<Vec<String>, crabidy_core::ProviderError> {
debug!("get_urls_for_track {}", track_uuid);
let (_, track_uuid, _) = split_uuid(track_uuid); let (_, track_uuid, _) = split_uuid(track_uuid);
let Ok(playback) = self.get_track_playback(&track_uuid).await else { let playback = self.get_track_playback(&track_uuid).await.map_err(|err| {
return Err(crabidy_core::ProviderError::FetchError); warn!(track = track_uuid, "failed to fetch playback info: {err}");
}; crabidy_core::ProviderError::FetchError
debug!("playback {:?}", playback); })?;
let Ok(manifest) = playback.get_manifest() else { trace!(?playback, "got playback info");
return Err(crabidy_core::ProviderError::FetchError); let manifest = playback.get_manifest().map_err(|err| {
}; warn!(track = track_uuid, "failed to decode manifest: {err}");
debug!("manifest {:?}", manifest); crabidy_core::ProviderError::FetchError
})?;
debug!(
track = track_uuid,
urls = manifest.urls.len(),
"resolved stream urls"
);
Ok(manifest.urls) Ok(manifest.urls)
} }
@ -66,10 +67,10 @@ impl crabidy_core::ProviderClient for Client {
&self, &self,
track_uuid: &str, track_uuid: &str,
) -> Result<crabidy_core::proto::crabidy::Track, crabidy_core::ProviderError> { ) -> Result<crabidy_core::proto::crabidy::Track, crabidy_core::ProviderError> {
debug!("get_metadata_for_track {}", track_uuid); let track = self.get_track(track_uuid).await.map_err(|err| {
let Ok(track) = self.get_track(track_uuid).await else { warn!(track = track_uuid, "failed to fetch track metadata: {err}");
return Err(crabidy_core::ProviderError::FetchError); crabidy_core::ProviderError::FetchError
}; })?;
Ok(track.into()) Ok(track.into())
} }
@ -107,9 +108,8 @@ impl crabidy_core::ProviderClient for Client {
let Some(user_id) = self.settings.login.user_id.clone() else { let Some(user_id) = self.settings.login.user_id.clone() else {
return Err(crabidy_core::ProviderError::UnknownUser); return Err(crabidy_core::ProviderError::UnknownUser);
}; };
debug!("get_lib_node in tidaldy{}", uuid);
let (_kind, module, uuid) = split_uuid(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() { let node = match module.as_str() {
"userplaylists" => { "userplaylists" => {
let mut node = crabidy_core::proto::crabidy::LibraryNode { let mut node = crabidy_core::proto::crabidy::LibraryNode {
@ -166,7 +166,6 @@ impl crabidy_core::ProviderClient for Client {
node node
} }
"artist" => { "artist" => {
info!("artist");
let mut node: crabidy_core::proto::crabidy::LibraryNode = let mut node: crabidy_core::proto::crabidy::LibraryNode =
self.get_artist(&uuid).await?.into(); self.get_artist(&uuid).await?.into();
let children: Vec<crabidy_core::proto::crabidy::LibraryNodeChild> = self let children: Vec<crabidy_core::proto::crabidy::LibraryNodeChild> = self
@ -371,7 +370,7 @@ impl Client {
error!("{:?}", e); error!("{:?}", e);
e e
})?; })?;
println!("{:?}", response); debug!(?response, "explorer response");
Ok(()) Ok(())
} }
@ -504,7 +503,13 @@ impl Client {
pub async fn login_web(&mut self) -> Result<(), ClientError> { pub async fn login_web(&mut self) -> Result<(), ClientError> {
let code_response = self.get_device_code().await?; let code_response = self.get_device_code().await?;
let now = Instant::now(); let now = Instant::now();
// The verification link must reach the user even without a log
// subscriber configured.
println!("https://{}", code_response.verification_uri_complete); 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 { while now.elapsed().as_secs() <= code_response.expires_in {
let login = self.check_auth_status(&code_response.device_code).await; let login = self.check_auth_status(&code_response.device_code).await;
if login.is_err() { if login.is_err() {
@ -520,9 +525,10 @@ impl Client {
self.settings.login.expires_after = Some(login_results.expires_in + timestamp); 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.user_id = Some(login_results.user.user_id.to_string());
self.settings.login.country_code = Some(login_results.user.country_code); self.settings.login.country_code = Some(login_results.user.country_code);
info!("device login succeeded");
return Ok(()); return Ok(());
} }
println!("login attempt expired"); warn!("device login attempt expired");
Err(ClientError::ConnectionError) Err(ClientError::ConnectionError)
} }