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, ClearQueueResponse, GetLibraryNodeRequest, GetLibraryNodeResponse, GetUpdateStreamRequest, GetUpdateStreamResponse, InitRequest, InitResponse, InsertRequest, InsertResponse, NextRequest, NextResponse, PrevRequest, PrevResponse, QueueRequest, QueueResponse, RemoveRequest, RemoveResponse, ReplaceRequest, ReplaceResponse, RestartTrackRequest, RestartTrackResponse, SaveQueueRequest, SaveQueueResponse, SetCurrentRequest, SetCurrentResponse, StopRequest, StopResponse, ToggleMuteRequest, ToggleMuteResponse, TogglePlayRequest, TogglePlayResponse, ToggleRepeatRequest, ToggleRepeatResponse, ToggleShuffleRequest, ToggleShuffleResponse, }; use std::pin::Pin; use tokio_stream::StreamExt; use tonic::{Request, Response, Status}; use tracing::{debug, error, instrument, trace}; #[derive(Debug)] pub struct RpcService { update_tx: tokio::sync::broadcast::Sender, playback_tx: flume::Sender, provider_tx: flume::Sender, } impl RpcService { pub fn new( update_tx: tokio::sync::broadcast::Sender, playback_tx: flume::Sender, provider_tx: flume::Sender, ) -> Self { Self { 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] impl CrabidyService for RpcService { type GetUpdateStreamStream = Pin> + Send>>; #[instrument(skip(self, _request))] async fn init(&self, _request: Request) -> Result, Status> { debug!("received init request"); let (result_tx, result_rx) = flume::bounded(1); 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)) } #[instrument(skip(self, request), fields(path))] async fn get_library_node( &self, request: Request, ) -> Result, Status> { let path = request.into_inner().path; tracing::Span::current().record("path", path.as_str()); debug!("received get_library_node request"); let (result_tx, result_rx) = flume::bounded(1); self.provider_tx .send_async(ProviderMessage::new(ProviderCommand::GetLibraryNode { path, result_tx, })) .await .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!("get_library_node failed: {err}"); Err(Status::internal(err.to_string())) } } } #[instrument(skip(self, request), fields(paths))] async fn queue( &self, request: Request, ) -> Result, Status> { let paths = request.into_inner().paths; tracing::Span::current().record("paths", format!("{paths:?}")); debug!("received queue request"); self.send_playback(PlaybackCommand::Queue { paths }).await?; Ok(Response::new(QueueResponse {})) } #[instrument(skip(self, request), fields(paths))] async fn replace( &self, request: Request, ) -> Result, Status> { let paths = request.into_inner().paths; tracing::Span::current().record("paths", format!("{paths:?}")); debug!("received replace request"); self.send_playback(PlaybackCommand::Replace { paths }) .await?; Ok(Response::new(ReplaceResponse {})) } #[instrument(skip(self, request), fields(paths))] async fn append( &self, request: Request, ) -> Result, Status> { let paths = request.into_inner().paths; tracing::Span::current().record("paths", format!("{paths:?}")); debug!("received append request"); self.send_playback(PlaybackCommand::Append { paths }) .await?; Ok(Response::new(AppendResponse {})) } #[instrument(skip(self, request), fields(positions))] async fn remove( &self, request: Request, ) -> Result, Status> { let positions = request.into_inner().positions; 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(paths, position))] async fn insert( &self, request: Request, ) -> Result, Status> { let req = request.into_inner(); tracing::Span::current().record("paths", format!("{:?}", req.paths)); tracing::Span::current().record("position", req.position); debug!("received insert request"); self.send_playback(PlaybackCommand::Insert { position: req.position, paths: req.paths, }) .await?; Ok(Response::new(InsertResponse {})) } #[instrument(skip(self, request), fields(exclude_current))] async fn clear_queue( &self, request: Request, ) -> Result, Status> { let exclude_current = request.into_inner().exclude_current; 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: Request, ) -> Result, Status> { let position = request.into_inner().position; 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: 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: 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: 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.map(|update_result| { trace!(?update_result, "forwarding update"); match update_result { Ok(update) => Ok(GetUpdateStreamResponse { update: Some(update), }), Err(err) => { // The client lagged too far behind the broadcast channel. error!("update stream lagged: {err}"); Err(Status::data_loss("update stream lagged")) } } }); Ok(Response::new(Box::pin(output_stream))) } #[instrument(skip(self, _request))] async fn save_queue( &self, _request: Request, ) -> Result, Status> { debug!("received save_queue request (not implemented)"); Ok(Response::new(SaveQueueResponse {})) } #[instrument(skip(self, _request))] async fn toggle_play( &self, _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: 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: Request, ) -> Result, Status> { let delta = request.into_inner().delta; 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: 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: 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: 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: Request, ) -> Result, Status> { debug!("received restart_track request"); self.send_playback(PlaybackCommand::RestartTrack).await?; Ok(Response::new(RestartTrackResponse {})) } }