feat: add metadata-agregator gRPC client integration

- Add tonic/prost for gRPC client generation
- Add proto definitions from metadata-agregator service
- Add MetadataClient and MetadataService for gRPC communication
- Add REST controller exposing metadata endpoints (search, get, sync)
- Update config with metadata.endpoint setting
- Update flake.nix with protobuf for proto compilation
This commit is contained in:
Alexander
2026-04-28 18:58:31 +02:00
parent 1aaaab4640
commit 5afcbd68ad
15 changed files with 1250 additions and 21 deletions
+342
View File
@@ -0,0 +1,342 @@
use axum::{
extract::{Path, Query, State},
http::StatusCode,
routing::{get, post},
Json, Router,
};
use serde::{Deserialize, Serialize};
use crate::metadata::proto::{Album, Artist, ArtistCredit, ExternalId, Genre, Label, Track, Work};
use crate::AppState;
pub fn routes() -> Router<AppState> {
Router::new()
.route("/artists/search", get(search_artists))
.route("/artists/:id", get(get_artist))
.route("/artists/:id/albums", get(get_artist_albums))
.route("/artists/sync", post(sync_artist))
.route("/albums/:id", get(get_album))
.route("/albums/:id/tracks", get(get_album_tracks))
.route("/status", get(connection_status))
}
#[derive(Debug, Deserialize)]
pub struct SearchQuery {
pub q: String,
pub limit: Option<i32>,
pub offset: Option<i32>,
}
#[derive(Debug, Serialize)]
pub struct ArtistResponse {
pub id: String,
pub name: String,
pub sort_name: String,
pub artist_type: String,
pub country: String,
pub formed_date: String,
pub disbanded_date: String,
pub description: String,
pub image_url: String,
pub genres: Vec<GenreResponse>,
pub external_ids: Vec<ExternalIdResponse>,
}
#[derive(Debug, Serialize)]
pub struct GenreResponse {
pub id: String,
pub name: String,
}
#[derive(Debug, Serialize)]
pub struct ExternalIdResponse {
pub source: String,
pub source_id: String,
pub url: String,
}
#[derive(Debug, Serialize)]
pub struct AlbumResponse {
pub id: String,
pub title: String,
pub album_type: String,
pub release_date: String,
pub upc: String,
pub total_tracks: i32,
pub total_discs: i32,
pub cover_url: String,
pub artists: Vec<ArtistCreditResponse>,
pub label: Option<LabelResponse>,
pub genres: Vec<GenreResponse>,
pub external_ids: Vec<ExternalIdResponse>,
}
#[derive(Debug, Serialize)]
pub struct ArtistCreditResponse {
pub artist: Option<ArtistResponse>,
pub role: String,
pub position: i32,
pub join_phrase: String,
}
#[derive(Debug, Serialize)]
pub struct LabelResponse {
pub id: String,
pub name: String,
pub country: String,
}
#[derive(Debug, Serialize)]
pub struct TrackResponse {
pub id: String,
pub title: String,
pub duration_ms: i32,
pub isrc: String,
pub explicit: bool,
pub disc_number: i32,
pub track_number: i32,
pub artists: Vec<ArtistCreditResponse>,
pub work: Option<WorkResponse>,
pub external_ids: Vec<ExternalIdResponse>,
}
#[derive(Debug, Serialize)]
pub struct WorkResponse {
pub id: String,
pub title: String,
pub work_type: String,
pub language: String,
}
#[derive(Debug, Serialize)]
pub struct SearchArtistsResponse {
pub artists: Vec<ArtistResponse>,
pub total: i32,
}
#[derive(Debug, Serialize)]
pub struct ArtistAlbumsResponse {
pub albums: Vec<AlbumResponse>,
pub total: i32,
}
#[derive(Debug, Serialize)]
pub struct AlbumTracksResponse {
pub tracks: Vec<TrackResponse>,
}
#[derive(Debug, Deserialize)]
pub struct SyncRequest {
pub name: String,
}
#[derive(Debug, Serialize)]
pub struct SyncResponse {
pub artist: Option<ArtistResponse>,
pub albums_synced: i32,
pub tracks_synced: i32,
}
fn map_genre(g: &Genre) -> GenreResponse {
GenreResponse {
id: g.id.clone(),
name: g.name.clone(),
}
}
fn map_external_id(e: &ExternalId) -> ExternalIdResponse {
ExternalIdResponse {
source: e.source.clone(),
source_id: e.source_id.clone(),
url: e.url.clone(),
}
}
fn map_artist(a: &Artist) -> ArtistResponse {
ArtistResponse {
id: a.id.clone(),
name: a.name.clone(),
sort_name: a.sort_name.clone(),
artist_type: a.artist_type.clone(),
country: a.country.clone(),
formed_date: a.formed_date.clone(),
disbanded_date: a.disbanded_date.clone(),
description: a.description.clone(),
image_url: a.image_url.clone(),
genres: a.genres.iter().map(map_genre).collect(),
external_ids: a.external_ids.iter().map(map_external_id).collect(),
}
}
fn map_label(l: &Label) -> LabelResponse {
LabelResponse {
id: l.id.clone(),
name: l.name.clone(),
country: l.country.clone(),
}
}
fn map_artist_credit(c: &ArtistCredit) -> ArtistCreditResponse {
ArtistCreditResponse {
artist: c.artist.as_ref().map(map_artist),
role: c.role.clone(),
position: c.position,
join_phrase: c.join_phrase.clone(),
}
}
fn map_album(a: &Album) -> AlbumResponse {
AlbumResponse {
id: a.id.clone(),
title: a.title.clone(),
album_type: a.album_type.clone(),
release_date: a.release_date.clone(),
upc: a.upc.clone(),
total_tracks: a.total_tracks,
total_discs: a.total_discs,
cover_url: a.cover_url.clone(),
artists: a.artists.iter().map(map_artist_credit).collect(),
label: a.label.as_ref().map(map_label),
genres: a.genres.iter().map(map_genre).collect(),
external_ids: a.external_ids.iter().map(map_external_id).collect(),
}
}
fn map_work(w: &Work) -> WorkResponse {
WorkResponse {
id: w.id.clone(),
title: w.title.clone(),
work_type: w.work_type.clone(),
language: w.language.clone(),
}
}
fn map_track(t: &Track) -> TrackResponse {
TrackResponse {
id: t.id.clone(),
title: t.title.clone(),
duration_ms: t.duration_ms,
isrc: t.isrc.clone(),
explicit: t.explicit,
disc_number: t.disc_number,
track_number: t.track_number,
artists: t.artists.iter().map(map_artist_credit).collect(),
work: t.work.as_ref().map(map_work),
external_ids: t.external_ids.iter().map(map_external_id).collect(),
}
}
async fn search_artists(
State(state): State<AppState>,
Query(query): Query<SearchQuery>,
) -> Result<Json<SearchArtistsResponse>, (StatusCode, String)> {
let state = state.read().await;
let response = state
.metadata_service
.search_artists(&query.q, query.limit, query.offset)
.await
.map_err(|e| (StatusCode::BAD_GATEWAY, e.to_string()))?;
Ok(Json(SearchArtistsResponse {
artists: response.artists.iter().map(map_artist).collect(),
total: response.total,
}))
}
async fn get_artist(
State(state): State<AppState>,
Path(id): Path<String>,
) -> Result<Json<ArtistResponse>, (StatusCode, String)> {
let state = state.read().await;
let artist = state
.metadata_service
.get_artist(&id)
.await
.map_err(|e| (StatusCode::BAD_GATEWAY, e.to_string()))?;
Ok(Json(map_artist(&artist)))
}
async fn get_artist_albums(
State(state): State<AppState>,
Path(id): Path<String>,
Query(query): Query<PaginationQuery>,
) -> Result<Json<ArtistAlbumsResponse>, (StatusCode, String)> {
let state = state.read().await;
let response = state
.metadata_service
.get_artist_albums(&id, query.limit, query.offset)
.await
.map_err(|e| (StatusCode::BAD_GATEWAY, e.to_string()))?;
Ok(Json(ArtistAlbumsResponse {
albums: response.albums.iter().map(map_album).collect(),
total: response.total,
}))
}
#[derive(Debug, Deserialize)]
pub struct PaginationQuery {
pub limit: Option<i32>,
pub offset: Option<i32>,
}
async fn get_album(
State(state): State<AppState>,
Path(id): Path<String>,
) -> Result<Json<AlbumResponse>, (StatusCode, String)> {
let state = state.read().await;
let album = state
.metadata_service
.get_album(&id)
.await
.map_err(|e| (StatusCode::BAD_GATEWAY, e.to_string()))?;
Ok(Json(map_album(&album)))
}
async fn get_album_tracks(
State(state): State<AppState>,
Path(id): Path<String>,
) -> Result<Json<AlbumTracksResponse>, (StatusCode, String)> {
let state = state.read().await;
let response = state
.metadata_service
.get_album_tracks(&id)
.await
.map_err(|e| (StatusCode::BAD_GATEWAY, e.to_string()))?;
Ok(Json(AlbumTracksResponse {
tracks: response.tracks.iter().map(map_track).collect(),
}))
}
async fn sync_artist(
State(state): State<AppState>,
Json(req): Json<SyncRequest>,
) -> Result<Json<SyncResponse>, (StatusCode, String)> {
let state = state.read().await;
let response = state
.metadata_service
.sync_artist(&req.name)
.await
.map_err(|e| (StatusCode::BAD_GATEWAY, e.to_string()))?;
Ok(Json(SyncResponse {
artist: response.artist.as_ref().map(map_artist),
albums_synced: response.albums_synced,
tracks_synced: response.tracks_synced,
}))
}
#[derive(Debug, Serialize)]
pub struct StatusResponse {
pub connected: bool,
}
async fn connection_status(State(state): State<AppState>) -> Json<StatusResponse> {
let state = state.read().await;
Json(StatusResponse {
connected: state.metadata_service.is_connected(),
})
}
+2
View File
@@ -1,4 +1,5 @@
mod indexer_controller;
mod metadata_controller;
mod torrent_controller;
use axum::{
@@ -23,6 +24,7 @@ pub fn routes(state: AppState) -> Router {
.route("/stats", get(get_stats))
.nest("/indexers", indexer_controller::routes())
.nest("/torrents", torrent_controller::routes())
.nest("/metadata", metadata_controller::routes())
.with_state(state)
}
+6
View File
@@ -15,10 +15,16 @@ pub enum ConfigError {
#[derive(Debug, Clone, Deserialize)]
pub struct Config {
pub database: DatabaseConfig,
pub metadata: MetadataConfig,
pub indexers: Vec<IndexerConfig>,
pub torrent: TorrentConfig,
}
#[derive(Debug, Clone, Deserialize)]
pub struct MetadataConfig {
pub endpoint: String,
}
#[derive(Debug, Clone, Deserialize)]
pub struct DatabaseConfig {
pub url: String,
+4
View File
@@ -1,6 +1,7 @@
pub mod api;
pub mod config;
pub mod indexer;
pub mod metadata;
pub mod models;
pub mod services;
pub mod torrent;
@@ -12,17 +13,20 @@ pub struct AppServices {
pub aggregator: services::Aggregator,
pub indexer_service: services::IndexerService,
pub torrent_service: services::TorrentService,
pub metadata_service: services::MetadataService,
}
impl AppServices {
pub fn new(
indexer_service: services::IndexerService,
torrent_service: services::TorrentService,
metadata_service: services::MetadataService,
) -> Self {
Self {
aggregator: services::Aggregator::new(),
indexer_service,
torrent_service,
metadata_service,
}
}
}
+18 -1
View File
@@ -4,7 +4,7 @@ use tokio::sync::RwLock;
use axum::Router;
use music_agregator::{
api, config,
services::{IndexerService, TorrentService},
services::{IndexerService, MetadataService, TorrentService},
AppServices, AppState,
};
use tower_http::cors::{Any, CorsLayer};
@@ -60,9 +60,26 @@ async fn main() {
TorrentService::new()
};
let mut metadata_service = MetadataService::new(&config.metadata.endpoint);
match metadata_service.connect().await {
Ok(()) => {
tracing::info!(
"connected to metadata service at {}",
config.metadata.endpoint
);
}
Err(e) => {
tracing::warn!(
"failed to connect to metadata service: {} (continuing without metadata)",
e
);
}
}
let state: AppState = Arc::new(RwLock::new(AppServices::new(
indexer_service,
torrent_service,
metadata_service,
)));
let cors = CorsLayer::new()
+112
View File
@@ -0,0 +1,112 @@
use thiserror::Error;
use tonic::transport::Channel;
use super::proto::{
get_album_request, get_artist_request, metadata_service_client::MetadataServiceClient,
sync_artist_request, Album, Artist, GetAlbumRequest, GetAlbumTracksRequest,
GetAlbumTracksResponse, GetArtistAlbumsRequest, GetArtistAlbumsResponse, GetArtistRequest,
Provider, SearchArtistsRequest, SearchArtistsResponse, SyncArtistRequest, SyncArtistResponse,
};
#[derive(Debug, Error)]
pub enum MetadataClientError {
#[error("connection failed: {0}")]
ConnectionFailed(String),
#[error("request failed: {0}")]
RequestFailed(#[from] tonic::Status),
#[error("transport error: {0}")]
Transport(#[from] tonic::transport::Error),
}
pub struct MetadataClient {
client: MetadataServiceClient<Channel>,
}
impl MetadataClient {
pub async fn connect(endpoint: &str) -> Result<Self, MetadataClientError> {
let client = MetadataServiceClient::connect(endpoint.to_string()).await?;
Ok(Self { client })
}
pub async fn search_artists(
&mut self,
query: &str,
limit: Option<i32>,
offset: Option<i32>,
) -> Result<SearchArtistsResponse, MetadataClientError> {
let request = SearchArtistsRequest {
query: query.to_string(),
limit: limit.unwrap_or(20),
offset: offset.unwrap_or(0),
provider: Provider::Unspecified as i32,
};
let response = self.client.search_artists(request).await?;
Ok(response.into_inner())
}
pub async fn get_artist(&mut self, id: &str) -> Result<Artist, MetadataClientError> {
let request = GetArtistRequest {
identifier: Some(get_artist_request::Identifier::Id(id.to_string())),
provider: Provider::Unspecified as i32,
};
let response = self.client.get_artist(request).await?;
Ok(response.into_inner())
}
pub async fn get_artist_albums(
&mut self,
artist_id: &str,
limit: Option<i32>,
offset: Option<i32>,
) -> Result<GetArtistAlbumsResponse, MetadataClientError> {
let request = GetArtistAlbumsRequest {
artist_id: artist_id.to_string(),
limit: limit.unwrap_or(50),
offset: offset.unwrap_or(0),
provider: Provider::Unspecified as i32,
};
let response = self.client.get_artist_albums(request).await?;
Ok(response.into_inner())
}
pub async fn get_album(&mut self, id: &str) -> Result<Album, MetadataClientError> {
let request = GetAlbumRequest {
identifier: Some(get_album_request::Identifier::Id(id.to_string())),
provider: Provider::Unspecified as i32,
};
let response = self.client.get_album(request).await?;
Ok(response.into_inner())
}
pub async fn get_album_tracks(
&mut self,
album_id: &str,
) -> Result<GetAlbumTracksResponse, MetadataClientError> {
let request = GetAlbumTracksRequest {
album_id: album_id.to_string(),
provider: Provider::Unspecified as i32,
};
let response = self.client.get_album_tracks(request).await?;
Ok(response.into_inner())
}
pub async fn sync_artist(
&mut self,
name: &str,
) -> Result<SyncArtistResponse, MetadataClientError> {
let request = SyncArtistRequest {
target: Some(sync_artist_request::Target::Name(name.to_string())),
provider: Provider::Musicbrainz as i32,
};
let response = self.client.sync_artist(request).await?;
Ok(response.into_inner())
}
}
+7
View File
@@ -0,0 +1,7 @@
mod client;
pub use client::{MetadataClient, MetadataClientError};
pub mod proto {
tonic::include_proto!("metadata.v1");
}
+92
View File
@@ -0,0 +1,92 @@
use std::sync::Arc;
use tokio::sync::Mutex;
use crate::metadata::{MetadataClient, MetadataClientError};
pub struct MetadataService {
client: Option<Arc<Mutex<MetadataClient>>>,
endpoint: String,
}
impl MetadataService {
pub fn new(endpoint: &str) -> Self {
Self {
client: None,
endpoint: endpoint.to_string(),
}
}
pub async fn connect(&mut self) -> Result<(), MetadataClientError> {
let client = MetadataClient::connect(&self.endpoint).await?;
self.client = Some(Arc::new(Mutex::new(client)));
Ok(())
}
pub fn is_connected(&self) -> bool {
self.client.is_some()
}
fn client(&self) -> Result<Arc<Mutex<MetadataClient>>, MetadataClientError> {
self.client
.clone()
.ok_or_else(|| MetadataClientError::ConnectionFailed("not connected".into()))
}
pub async fn search_artists(
&self,
query: &str,
limit: Option<i32>,
offset: Option<i32>,
) -> Result<crate::metadata::proto::SearchArtistsResponse, MetadataClientError> {
let client = self.client()?;
let mut guard = client.lock().await;
guard.search_artists(query, limit, offset).await
}
pub async fn get_artist(
&self,
id: &str,
) -> Result<crate::metadata::proto::Artist, MetadataClientError> {
let client = self.client()?;
let mut guard = client.lock().await;
guard.get_artist(id).await
}
pub async fn get_artist_albums(
&self,
artist_id: &str,
limit: Option<i32>,
offset: Option<i32>,
) -> Result<crate::metadata::proto::GetArtistAlbumsResponse, MetadataClientError> {
let client = self.client()?;
let mut guard = client.lock().await;
guard.get_artist_albums(artist_id, limit, offset).await
}
pub async fn get_album(
&self,
id: &str,
) -> Result<crate::metadata::proto::Album, MetadataClientError> {
let client = self.client()?;
let mut guard = client.lock().await;
guard.get_album(id).await
}
pub async fn get_album_tracks(
&self,
album_id: &str,
) -> Result<crate::metadata::proto::GetAlbumTracksResponse, MetadataClientError> {
let client = self.client()?;
let mut guard = client.lock().await;
guard.get_album_tracks(album_id).await
}
pub async fn sync_artist(
&self,
name: &str,
) -> Result<crate::metadata::proto::SyncArtistResponse, MetadataClientError> {
let client = self.client()?;
let mut guard = client.lock().await;
guard.sync_artist(name).await
}
}
+2
View File
@@ -1,7 +1,9 @@
mod indexer_service;
mod metadata_service;
mod torrent_service;
pub use indexer_service::{IndexerInfo, IndexerService};
pub use metadata_service::MetadataService;
pub use torrent_service::TorrentService;
use uuid::Uuid;