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, debug_span, error, instrument, warn, Instrument}; #[derive(Debug)] pub struct ProviderOrchestrator { pub provider_tx: flume::Sender, provider_rx: flume::Receiver, tidal_client: Arc, } impl ProviderOrchestrator { pub fn run(self) { tokio::spawn(async move { 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 { path, result_tx } => { let result = self.get_lib_node(&path).await; if let Err(err) = result_tx.send_async(result).await { error!("failed to send get_library_node result: {err}"); } } ProviderCommand::GetTrackUrls { path, result_tx } => { let result = self.get_urls_for_track(&path).await; if let Err(err) = result_tx.send_async(result).await { error!("failed to send get_track_urls result: {err}"); } } ProviderCommand::ResolveTracks { path, result_tx } => { let result = self.resolve_tracks(&path).await; if let Err(err) = result_tx.send_async(result).await { error!("failed to send resolve_tracks result: {err}"); } } } } /// Resolves a path into playable tracks. A track path resolves to that /// single track; a node path is flattened by walking its queueable /// descendants. #[instrument(skip(self))] async fn resolve_tracks(&self, path: &str) -> Vec { if self.is_track_path(path) { return match self.get_metadata_for_track(path).await { Ok(track) => vec![track], Err(err) => { warn!(path, "failed to resolve track: {err}"); Vec::new() } }; } let mut tracks = Vec::new(); let mut nodes_to_go = vec![path.to_string()]; while let Some(node_path) = nodes_to_go.pop() { let node = match self.get_lib_node(&node_path).await { Ok(node) => node, Err(err) => { warn!(node = node_path, "skipping unreadable node: {err}"); continue; } }; if node.is_queable { tracks.extend(node.tracks); nodes_to_go.extend(node.children.into_iter().map(|c| c.path)) } } debug!(count = tracks.len(), "resolved path into tracks"); tracks } } #[async_trait] impl ProviderClient for ProviderOrchestrator { #[instrument(skip(_s))] async fn init(_s: &str) -> Result { let config_dir = dirs::config_dir() .map(|d| d.join("crabidy")) .unwrap_or(PathBuf::from("/tmp")); let dir_exists = tokio::fs::try_exists(&config_dir) .await .map_err(|e| ProviderError::Config(e.to_string()))?; if !dir_exists { tokio::fs::create_dir(&config_dir) .await .map_err(|e| ProviderError::Config(e.to_string()))?; } let config_file = config_dir.join("tidaly.toml"); 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).await { error!("failed to write tidal config file: {err}"); }; let (provider_tx, provider_rx) = flume::bounded(100); Ok(Self { provider_rx, provider_tx, tidal_client, }) } fn settings(&self) -> String { String::new() } /// Routes to the provider that owns the path. fn is_track_path(&self, path: &str) -> bool { if path == "/tidal" || path.starts_with("/tidal/") { return self.tidal_client.is_track_path(path); } false } #[instrument(skip(self))] async fn get_urls_for_track(&self, track_path: &str) -> Result, ProviderError> { if track_path.starts_with("/tidal/") { return self.tidal_client.get_urls_for_track(track_path).await; } warn!(path = track_path, "no provider owns this track path"); Err(ProviderError::MalformedPath) } #[instrument(skip(self))] async fn get_metadata_for_track(&self, track_path: &str) -> Result { if track_path.starts_with("/tidal/") { return self.tidal_client.get_metadata_for_track(track_path).await; } warn!(path = track_path, "no provider owns this track path"); Err(ProviderError::MalformedPath) } fn get_lib_root(&self) -> LibraryNode { let mut root_node = LibraryNode::new(); let child = LibraryNodeChild::new(tidaldy::PROVIDER_ROOT.to_owned(), "tidal".to_owned(), false); root_node.children.push(child); root_node } #[instrument(skip(self))] async fn get_lib_node(&self, path: &str) -> Result { if path == crabidy_core::ROOT_PATH { debug!("serving global library root"); return Ok(self.get_lib_root()); } if path == tidaldy::PROVIDER_ROOT || path.starts_with("/tidal/") { return self.tidal_client.get_lib_node(path).await; } warn!(path, "no provider owns this path"); Err(ProviderError::MalformedPath) } }