Add metadata provider
This commit is contained in:
@@ -1 +1 @@
|
||||
/nix/store/gd69vhx5mbadafs4bwhh66b5mnp3wvyg-pre-commit-config.json
|
||||
/nix/store/caqfpcm0y8ar1fh9apqimnndhq32cfc1-pre-commit-config.json
|
||||
@@ -0,0 +1,14 @@
|
||||
info:
|
||||
name: SearchArtists
|
||||
type: grpc
|
||||
seq: 1
|
||||
|
||||
grpc:
|
||||
url: localhost:50051
|
||||
method: /metadata.v1.MetadataService/SearchArtists
|
||||
methodType: unary
|
||||
message: |-
|
||||
{
|
||||
"query": "ДДТ"
|
||||
}
|
||||
auth: inherit
|
||||
@@ -0,0 +1,7 @@
|
||||
info:
|
||||
name: Metadata
|
||||
type: folder
|
||||
seq: 4
|
||||
|
||||
request:
|
||||
auth: inherit
|
||||
@@ -15,3 +15,4 @@ inputs:
|
||||
imports:
|
||||
- ./musicfs
|
||||
- ./tora
|
||||
- ./metadata-agregator
|
||||
|
||||
+1
-1
Submodule metadata-agregator updated: 995397e495...9f815f69c9
@@ -5,6 +5,7 @@ pub struct Config {
|
||||
pub service: ServiceConfig,
|
||||
pub indexer: IndexerConfig,
|
||||
pub torrent: TorrentConfig,
|
||||
pub metadata: MetadataConfig,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
@@ -12,6 +13,7 @@ pub struct ServiceConfig {
|
||||
pub address: String,
|
||||
}
|
||||
|
||||
// TODO move the indexer config to torrent and make it as subconfig
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct IndexerConfig {
|
||||
pub address: String,
|
||||
@@ -22,3 +24,8 @@ pub struct IndexerConfig {
|
||||
pub struct TorrentConfig {
|
||||
pub address: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct MetadataConfig {
|
||||
pub address: String,
|
||||
}
|
||||
|
||||
+60
-9
@@ -10,19 +10,25 @@ use tonic_health::pb::{
|
||||
use crate::generated::health::{
|
||||
CheckRequest, CheckResponse, ServingStatus, SubserviceStatus, health_server::Health,
|
||||
};
|
||||
use crate::generated::metadata::metadata_service_server;
|
||||
use crate::generated::torrent::torrents_server;
|
||||
|
||||
const QUERY_TIMEOUT: Duration = Duration::from_secs(2);
|
||||
const TORA_SUBSERVICE_NAME: &str = "tora";
|
||||
const METADATA_SUBSERVICE_NAME: &str = "metadata-agregator";
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct HealthService {
|
||||
tora_channel: Channel,
|
||||
metadata_channel: Channel,
|
||||
}
|
||||
|
||||
impl HealthService {
|
||||
pub fn new(tora_channel: Channel) -> Self {
|
||||
Self { tora_channel }
|
||||
pub fn new(tora_channel: Channel, metadata_channel: Channel) -> Self {
|
||||
Self {
|
||||
tora_channel,
|
||||
metadata_channel,
|
||||
}
|
||||
}
|
||||
|
||||
async fn probe_tora(&self) -> Result<(), String> {
|
||||
@@ -41,6 +47,26 @@ impl HealthService {
|
||||
Err(_) => Err("torad health check timed out".to_string()),
|
||||
}
|
||||
}
|
||||
|
||||
async fn probe_metadata(&self) -> Result<(), String> {
|
||||
let mut client = HealthClient::new(self.metadata_channel.clone());
|
||||
let request = GrpcHealthCheckRequest {
|
||||
service: String::new(),
|
||||
};
|
||||
|
||||
match tokio::time::timeout(QUERY_TIMEOUT, client.check(request)).await {
|
||||
Ok(Ok(response)) => match GrpcServingStatus::try_from(response.into_inner().status) {
|
||||
Ok(GrpcServingStatus::Serving) => Ok(()),
|
||||
Ok(status) => Err(format!(
|
||||
"metadata-agregator reported {}",
|
||||
status.as_str_name()
|
||||
)),
|
||||
Err(_) => Err("metadata-agregator reported an unknown status".to_string()),
|
||||
},
|
||||
Ok(Err(status)) => Err(format!("metadata-agregator health check failed: {status}")),
|
||||
Err(_) => Err("metadata-agregator health check timed out".to_string()),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tonic::async_trait]
|
||||
@@ -50,7 +76,20 @@ impl GrpcHealth for HealthService {
|
||||
request: tonic::Request<GrpcHealthCheckRequest>,
|
||||
) -> Result<tonic::Response<GrpcHealthCheckResponse>, tonic::Status> {
|
||||
let status = match request.into_inner().service.as_str() {
|
||||
"" | torrents_server::SERVICE_NAME => match self.probe_tora().await {
|
||||
"" => {
|
||||
let tora_ok = self.probe_tora().await.is_ok();
|
||||
let metadata_ok = self.probe_metadata().await.is_ok();
|
||||
if tora_ok && metadata_ok {
|
||||
GrpcServingStatus::Serving
|
||||
} else {
|
||||
GrpcServingStatus::NotServing
|
||||
}
|
||||
}
|
||||
torrents_server::SERVICE_NAME => match self.probe_tora().await {
|
||||
Ok(()) => GrpcServingStatus::Serving,
|
||||
Err(_) => GrpcServingStatus::NotServing,
|
||||
},
|
||||
metadata_service_server::SERVICE_NAME => match self.probe_metadata().await {
|
||||
Ok(()) => GrpcServingStatus::Serving,
|
||||
Err(_) => GrpcServingStatus::NotServing,
|
||||
},
|
||||
@@ -86,16 +125,28 @@ impl Health for HealthService {
|
||||
&self,
|
||||
_request: tonic::Request<CheckRequest>,
|
||||
) -> Result<tonic::Response<CheckResponse>, tonic::Status> {
|
||||
let (status, message) = match self.probe_tora().await {
|
||||
let (tora_status, tora_message) = match self.probe_tora().await {
|
||||
Ok(()) => (ServingStatus::Serving, String::new()),
|
||||
Err(message) => (ServingStatus::NotServing, message),
|
||||
};
|
||||
|
||||
let subservices = vec![SubserviceStatus {
|
||||
name: TORA_SUBSERVICE_NAME.to_string(),
|
||||
status: status as i32,
|
||||
message,
|
||||
}];
|
||||
let (metadata_status, metadata_message) = match self.probe_metadata().await {
|
||||
Ok(()) => (ServingStatus::Serving, String::new()),
|
||||
Err(message) => (ServingStatus::NotServing, message),
|
||||
};
|
||||
|
||||
let subservices = vec![
|
||||
SubserviceStatus {
|
||||
name: TORA_SUBSERVICE_NAME.to_string(),
|
||||
status: tora_status as i32,
|
||||
message: tora_message,
|
||||
},
|
||||
SubserviceStatus {
|
||||
name: METADATA_SUBSERVICE_NAME.to_string(),
|
||||
status: metadata_status as i32,
|
||||
message: metadata_message,
|
||||
},
|
||||
];
|
||||
|
||||
let overall = if subservices
|
||||
.iter()
|
||||
|
||||
+12
-2
@@ -2,6 +2,7 @@ mod config;
|
||||
mod greeter;
|
||||
mod health;
|
||||
mod indexer;
|
||||
mod metadata;
|
||||
mod torrent;
|
||||
mod torrent_manager;
|
||||
|
||||
@@ -18,6 +19,10 @@ mod generated {
|
||||
include!("generated/torrent_manager/torrent_manager.rs");
|
||||
}
|
||||
|
||||
pub mod metadata {
|
||||
include!("generated/metadata/v1/metadata.v1.rs");
|
||||
}
|
||||
|
||||
pub mod health {
|
||||
include!("generated/health/health.rs");
|
||||
}
|
||||
@@ -27,12 +32,14 @@ use std::{fs, sync::Arc};
|
||||
|
||||
use generated::{
|
||||
health::health_server::HealthServer as AggregatorHealthServer,
|
||||
hello::greeter_server::GreeterServer, torrent::torrents_server::TorrentsServer,
|
||||
hello::greeter_server::GreeterServer, metadata::metadata_service_server::MetadataServiceServer,
|
||||
torrent::torrents_server::TorrentsServer,
|
||||
torrent_manager::torrent_manager_server::TorrentManagerServer,
|
||||
};
|
||||
use greeter::GreeterService;
|
||||
use health::HealthService;
|
||||
use indexer::Jackett;
|
||||
use metadata::MetadataService;
|
||||
use tonic_health::pb::health_server::HealthServer;
|
||||
use torrent::TorrentsService;
|
||||
use torrent_manager::TorrentMananagerService;
|
||||
@@ -48,11 +55,13 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
.register_encoded_file_descriptor_set(generated::torrent::FILE_DESCRIPTOR_SET)
|
||||
.register_encoded_file_descriptor_set(generated::health::FILE_DESCRIPTOR_SET)
|
||||
.register_encoded_file_descriptor_set(generated::torrent_manager::FILE_DESCRIPTOR_SET)
|
||||
.register_encoded_file_descriptor_set(generated::metadata::FILE_DESCRIPTOR_SET)
|
||||
.register_encoded_file_descriptor_set(tonic_health::pb::FILE_DESCRIPTOR_SET)
|
||||
.build_v1()?;
|
||||
|
||||
let torrents = TorrentsService::new(&config.torrent.address);
|
||||
let health = HealthService::new(torrents.channel());
|
||||
let metadata = MetadataService::new(&config.metadata.address);
|
||||
let health = HealthService::new(torrents.channel(), metadata.channel());
|
||||
let indexer = Arc::new(Jackett::new(&config.indexer));
|
||||
let torrent_manager = TorrentMananagerService::new(torrents.clone(), indexer);
|
||||
|
||||
@@ -63,6 +72,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
.add_service(GreeterServer::new(GreeterService::default()))
|
||||
.add_service(TorrentsServer::new(torrents))
|
||||
.add_service(TorrentManagerServer::new(torrent_manager))
|
||||
.add_service(MetadataServiceServer::new(metadata))
|
||||
.serve(addr)
|
||||
.await?;
|
||||
|
||||
|
||||
@@ -0,0 +1,98 @@
|
||||
use crate::generated::metadata::{
|
||||
metadata_service_client::MetadataServiceClient,
|
||||
metadata_service_server::MetadataService as Metadata, *,
|
||||
};
|
||||
use tonic::transport::{Channel, Endpoint};
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct MetadataService {
|
||||
channel: Channel,
|
||||
}
|
||||
|
||||
impl MetadataService {
|
||||
pub fn new(address: &str) -> Self {
|
||||
let channel = Endpoint::from_shared(address.to_string())
|
||||
.expect("metadata.address should be a valid endpoint")
|
||||
.connect_lazy();
|
||||
Self { channel }
|
||||
}
|
||||
|
||||
pub fn channel(&self) -> Channel {
|
||||
self.channel.clone()
|
||||
}
|
||||
}
|
||||
|
||||
#[tonic::async_trait]
|
||||
impl Metadata for MetadataService {
|
||||
async fn get_artist(
|
||||
&self,
|
||||
request: tonic::Request<GetArtistRequest>,
|
||||
) -> Result<tonic::Response<GetArtistResponse>, tonic::Status> {
|
||||
MetadataServiceClient::new(self.channel.clone())
|
||||
.get_artist(request.into_inner())
|
||||
.await
|
||||
}
|
||||
|
||||
async fn search_artists(
|
||||
&self,
|
||||
request: tonic::Request<SearchArtistsRequest>,
|
||||
) -> Result<tonic::Response<SearchArtistsResponse>, tonic::Status> {
|
||||
MetadataServiceClient::new(self.channel.clone())
|
||||
.search_artists(request.into_inner())
|
||||
.await
|
||||
}
|
||||
|
||||
async fn get_album(
|
||||
&self,
|
||||
request: tonic::Request<GetAlbumRequest>,
|
||||
) -> Result<tonic::Response<GetAlbumResponse>, tonic::Status> {
|
||||
MetadataServiceClient::new(self.channel.clone())
|
||||
.get_album(request.into_inner())
|
||||
.await
|
||||
}
|
||||
|
||||
async fn get_artist_albums(
|
||||
&self,
|
||||
request: tonic::Request<GetArtistAlbumsRequest>,
|
||||
) -> Result<tonic::Response<GetArtistAlbumsResponse>, tonic::Status> {
|
||||
MetadataServiceClient::new(self.channel.clone())
|
||||
.get_artist_albums(request.into_inner())
|
||||
.await
|
||||
}
|
||||
|
||||
async fn get_track(
|
||||
&self,
|
||||
request: tonic::Request<GetTrackRequest>,
|
||||
) -> Result<tonic::Response<GetTrackResponse>, tonic::Status> {
|
||||
MetadataServiceClient::new(self.channel.clone())
|
||||
.get_track(request.into_inner())
|
||||
.await
|
||||
}
|
||||
|
||||
async fn get_album_tracks(
|
||||
&self,
|
||||
request: tonic::Request<GetAlbumTracksRequest>,
|
||||
) -> Result<tonic::Response<GetAlbumTracksResponse>, tonic::Status> {
|
||||
MetadataServiceClient::new(self.channel.clone())
|
||||
.get_album_tracks(request.into_inner())
|
||||
.await
|
||||
}
|
||||
|
||||
async fn search_albums(
|
||||
&self,
|
||||
request: tonic::Request<SearchAlbumsRequest>,
|
||||
) -> Result<tonic::Response<SearchAlbumsResponse>, tonic::Status> {
|
||||
MetadataServiceClient::new(self.channel.clone())
|
||||
.search_albums(request.into_inner())
|
||||
.await
|
||||
}
|
||||
|
||||
async fn sync_artist(
|
||||
&self,
|
||||
request: tonic::Request<SyncArtistRequest>,
|
||||
) -> Result<tonic::Response<SyncArtistResponse>, tonic::Status> {
|
||||
MetadataServiceClient::new(self.channel.clone())
|
||||
.sync_artist(request.into_inner())
|
||||
.await
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user