crabidy/crabidy-server/src/rpc.rs

526 lines
21 KiB
Rust

use crate::{PlaybackCommand, PlaybackMessage, ProviderCommand, ProviderMessage};
use crabidy_core::proto::crabidy::{
crabidy_service_server::CrabidyService, get_update_stream_response::Update as StreamUpdate,
AppendRequest, AppendResponse, CaptureLibraryNodeRequest, CaptureLibraryNodeResponse,
ChangeVolumeRequest, ChangeVolumeResponse, ClearQueueRequest, ClearQueueResponse,
CreateLibraryNodeRequest, CreateLibraryNodeResponse, DeleteLibraryNodeRequest,
DeleteLibraryNodeResponse, GetLibraryNodeRequest, GetLibraryNodeResponse,
GetUpdateStreamRequest, GetUpdateStreamResponse, InitRequest, InitResponse, InsertRequest,
InsertResponse, NextRequest, NextResponse, PrevRequest, PrevResponse, QueueRequest,
QueueResponse, RemoveRequest, RemoveResponse, RenameLibraryNodeRequest,
RenameLibraryNodeResponse, ReplaceRequest, ReplaceResponse, RestartTrackRequest,
RestartTrackResponse, SaveQueueRequest, SaveQueueResponse, SetCurrentRequest,
SetCurrentResponse, StopRequest, StopResponse, ToggleMuteRequest, ToggleMuteResponse,
TogglePlayRequest, TogglePlayResponse, ToggleRepeatRequest, ToggleRepeatResponse,
ToggleShuffleRequest, ToggleShuffleResponse,
};
use crabidy_core::ProviderError;
use crabidy_server::bookmark_store::CaptureError;
use crabidy_server::queue_store::SaveQueueError;
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<StreamUpdate>,
playback_tx: flume::Sender<PlaybackMessage>,
provider_tx: flume::Sender<ProviderMessage>,
}
impl RpcService {
pub fn new(
update_tx: tokio::sync::broadcast::Sender<StreamUpdate>,
playback_tx: flume::Sender<PlaybackMessage>,
provider_tx: flume::Sender<ProviderMessage>,
) -> 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<Box<dyn tokio_stream::Stream<Item = Result<GetUpdateStreamResponse, Status>> + Send>>;
#[instrument(skip(self, _request))]
async fn init(&self, _request: Request<InitRequest>) -> Result<Response<InitResponse>, 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<GetLibraryNodeRequest>,
) -> Result<Response<GetLibraryNodeResponse>, 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()))
}
}
}
/// Creates a node under a creatable parent via the provider loop.
///
/// Error mapping is part of the contract: `NotSupported` →
/// `failed_precondition`, `InvalidInput` → `invalid_argument`, everything
/// else `internal`.
#[instrument(skip(self, request), fields(parent_path, title))]
async fn create_library_node(
&self,
request: Request<CreateLibraryNodeRequest>,
) -> Result<Response<CreateLibraryNodeResponse>, Status> {
let CreateLibraryNodeRequest { parent_path, title } = request.into_inner();
tracing::Span::current().record("parent_path", parent_path.as_str());
tracing::Span::current().record("title", title.as_str());
debug!("received create_library_node request");
let (result_tx, result_rx) = flume::bounded(1);
self.provider_tx
.send_async(ProviderMessage::new(ProviderCommand::CreateLibraryNode {
parent_path,
title,
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(CreateLibraryNodeResponse {
node: Some(node),
})),
Err(ProviderError::NotSupported) => Err(Status::failed_precondition(
"this node does not support creating children",
)),
Err(ProviderError::InvalidInput) => Err(Status::invalid_argument("invalid node title")),
Err(err) => {
error!("create_library_node failed: {err}");
Err(Status::internal(err.to_string()))
}
}
}
/// Renames an editable node via the provider loop. Same error mapping as
/// `create_library_node`: `NotSupported` → `failed_precondition`,
/// `InvalidInput` → `invalid_argument`, everything else `internal`.
#[instrument(skip(self, request), fields(path, new_title))]
async fn rename_library_node(
&self,
request: Request<RenameLibraryNodeRequest>,
) -> Result<Response<RenameLibraryNodeResponse>, Status> {
let RenameLibraryNodeRequest { path, new_title } = request.into_inner();
tracing::Span::current().record("path", path.as_str());
tracing::Span::current().record("new_title", new_title.as_str());
debug!("received rename_library_node request");
let (result_tx, result_rx) = flume::bounded(1);
self.provider_tx
.send_async(ProviderMessage::new(ProviderCommand::RenameLibraryNode {
path,
new_title,
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(RenameLibraryNodeResponse {
node: Some(node),
})),
Err(ProviderError::NotSupported) => {
Err(Status::failed_precondition("this node cannot be renamed"))
}
Err(ProviderError::InvalidInput) => Err(Status::invalid_argument("invalid node title")),
Err(err) => {
error!("rename_library_node failed: {err}");
Err(Status::internal(err.to_string()))
}
}
}
/// Deletes a deletable node via the provider loop and returns the
/// refreshed parent. Same error mapping as `create_library_node`.
#[instrument(skip(self, request), fields(path))]
async fn delete_library_node(
&self,
request: Request<DeleteLibraryNodeRequest>,
) -> Result<Response<DeleteLibraryNodeResponse>, Status> {
let DeleteLibraryNodeRequest { path } = request.into_inner();
tracing::Span::current().record("path", path.as_str());
debug!("received delete_library_node request");
let (result_tx, result_rx) = flume::bounded(1);
self.provider_tx
.send_async(ProviderMessage::new(ProviderCommand::DeleteLibraryNode {
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(parent) => Ok(Response::new(DeleteLibraryNodeResponse {
parent: Some(parent),
})),
Err(ProviderError::NotSupported) => {
Err(Status::failed_precondition("this node cannot be deleted"))
}
// No provider raises this for delete today; mapped anyway so the
// contract stays uniform across the node-mutation rpcs.
Err(ProviderError::InvalidInput) => Err(Status::invalid_argument("invalid node path")),
Err(err) => {
error!("delete_library_node failed: {err}");
Err(Status::internal(err.to_string()))
}
}
}
#[instrument(skip(self, request), fields(paths))]
async fn queue(
&self,
request: Request<QueueRequest>,
) -> Result<Response<QueueResponse>, 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<ReplaceRequest>,
) -> Result<Response<ReplaceResponse>, 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<AppendRequest>,
) -> Result<Response<AppendResponse>, 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<RemoveRequest>,
) -> Result<Response<RemoveResponse>, 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<InsertRequest>,
) -> Result<Response<InsertResponse>, 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<ClearQueueRequest>,
) -> Result<Response<ClearQueueResponse>, 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<SetCurrentRequest>,
) -> Result<Response<SetCurrentResponse>, 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<ToggleShuffleRequest>,
) -> Result<Response<ToggleShuffleResponse>, 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<ToggleRepeatRequest>,
) -> Result<Response<ToggleRepeatResponse>, 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<GetUpdateStreamRequest>,
) -> Result<Response<Self::GetUpdateStreamStream>, 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)))
}
/// Captures a queueable subtree as a bookmark via the provider loop
/// (structure-preserving snapshot under `/bookmarks/<name>`).
///
/// Error mapping is part of the contract: invalid name or
/// uncapturable source → `invalid_argument`; an over-cap subtree or
/// disabled bookmarks → `failed_precondition`; walk/write failures →
/// `internal`.
#[instrument(skip(self, request), fields(path, name))]
async fn capture_library_node(
&self,
request: Request<CaptureLibraryNodeRequest>,
) -> Result<Response<CaptureLibraryNodeResponse>, Status> {
let CaptureLibraryNodeRequest { path, name } = request.into_inner();
tracing::Span::current().record("path", path.as_str());
tracing::Span::current().record("name", name.as_str());
debug!("received capture_library_node request");
let (result_tx, result_rx) = flume::bounded(1);
self.provider_tx
.send_async(ProviderMessage::new(ProviderCommand::CaptureLibraryNode {
path,
name,
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(()) => Ok(Response::new(CaptureLibraryNodeResponse {})),
Err(err @ (CaptureError::InvalidName(_) | CaptureError::BadSource(_))) => {
Err(Status::invalid_argument(err.to_string()))
}
Err(err @ (CaptureError::TooLarge(_) | CaptureError::Disabled)) => {
Err(Status::failed_precondition(err.to_string()))
}
Err(err) => {
error!("capture_library_node failed: {err}");
Err(Status::internal("cannot capture the subtree"))
}
}
}
/// Saves the current queue under a name (persisted queues, visible as
/// `/queues/<name>` in the library).
///
/// Error mapping is part of the contract: an invalid name →
/// `invalid_argument`; an empty queue or disabled persistence →
/// `failed_precondition`; I/O and serialization failures → `internal`.
#[instrument(skip(self, request), fields(name))]
async fn save_queue(
&self,
request: Request<SaveQueueRequest>,
) -> Result<Response<SaveQueueResponse>, Status> {
let name = request.into_inner().name;
tracing::Span::current().record("name", name.as_str());
debug!("received save_queue request");
let (result_tx, result_rx) = flume::bounded(1);
self.send_playback(PlaybackCommand::SaveQueue { name, result_tx })
.await?;
let result = result_rx.recv_async().await.map_err(|err| {
error!("no reply from playback loop: {err}");
Status::internal("playback loop did not reply")
})?;
match result {
Ok(()) => Ok(Response::new(SaveQueueResponse {})),
Err(err @ SaveQueueError::InvalidName(_)) => {
Err(Status::invalid_argument(err.to_string()))
}
Err(err @ (SaveQueueError::EmptyQueue | SaveQueueError::Disabled)) => {
Err(Status::failed_precondition(err.to_string()))
}
Err(err) => {
error!("save_queue failed: {err}");
Err(Status::internal("cannot save the queue"))
}
}
}
#[instrument(skip(self, _request))]
async fn toggle_play(
&self,
_request: Request<TogglePlayRequest>,
) -> Result<Response<TogglePlayResponse>, 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<StopRequest>) -> Result<Response<StopResponse>, 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<ChangeVolumeRequest>,
) -> Result<Response<ChangeVolumeResponse>, 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<ToggleMuteRequest>,
) -> Result<Response<ToggleMuteResponse>, 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<NextRequest>) -> Result<Response<NextResponse>, 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<PrevRequest>) -> Result<Response<PrevResponse>, 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<RestartTrackRequest>,
) -> Result<Response<RestartTrackResponse>, Status> {
debug!("received restart_track request");
self.send_playback(PlaybackCommand::RestartTrack).await?;
Ok(Response::new(RestartTrackResponse {}))
}
}