From 8a2e379f6df3cab92d5822ac32a4c62facb6030a Mon Sep 17 00:00:00 2001 From: Juanjo Presa Date: Mon, 4 May 2026 20:22:51 +0200 Subject: [PATCH] feat: add media-server input foundations (#742) * feat: add remote media input foundations * fixed macro variable names * Added todo comment for source editor sidebar --- .../src/api/endpoints/api_playlist_utils.rs | 7 + .../src/api/model/metadata_update_manager.rs | 23 +- .../src/api/model/provider_lineup_manager.rs | 1 + backend/src/api/setup_api.rs | 2 +- backend/src/media_server/catalog.rs | 307 ++++++++ backend/src/media_server/client.rs | 128 ++++ backend/src/media_server/emby/dto.rs | 88 +++ backend/src/media_server/emby/mod.rs | 1 + backend/src/media_server/errors.rs | 190 +++++ backend/src/media_server/jellyfin/dto.rs | 26 + backend/src/media_server/jellyfin/mod.rs | 1 + backend/src/media_server/mod.rs | 27 + backend/src/media_server/playback.rs | 335 +++++++++ backend/src/media_server/playlist_mapper.rs | 249 ++++++ backend/src/media_server/plex/dto.rs | 198 +++++ backend/src/media_server/plex/mod.rs | 1 + backend/src/media_server/redaction.rs | 170 +++++ backend/src/media_server/test_fixtures.rs | 65 ++ backend/src/media_server/types.rs | 269 +++++++ backend/src/model/config/input.rs | 356 ++++++++- backend/src/modules.rs | 1 + backend/src/processing/processor/playlist.rs | 10 + .../src/processing/processor/stream_probe.rs | 12 +- backend/src/repository/alias_repository.rs | 1 + backend/src/repository/playlist_repository.rs | 58 +- backend/src/repository/playlist_source.rs | 1 + frontend/public/assets/i18n/en.json | 3 + frontend/src/app/components/mod.rs | 5 +- .../playlist/input/input_type_view.rs | 3 + .../playlist/input/staged_input_view.rs | 3 + .../components/source_editor/editor_model.rs | 46 +- .../components/source_editor/editor_view.rs | 54 ++ .../components/source_editor/output_form.rs | 8 +- .../app/components/source_editor/sidebar.rs | 10 +- frontend/src/services/config_service.rs | 5 +- shared/src/model/config/input.rs | 708 +++++++++++++++++- shared/src/model/config/macros.rs | 54 +- shared/src/model/xtream.rs | 4 +- shared/src/utils/default_utils.rs | 2 + 39 files changed, 3376 insertions(+), 56 deletions(-) create mode 100644 backend/src/media_server/catalog.rs create mode 100644 backend/src/media_server/client.rs create mode 100644 backend/src/media_server/emby/dto.rs create mode 100644 backend/src/media_server/emby/mod.rs create mode 100644 backend/src/media_server/errors.rs create mode 100644 backend/src/media_server/jellyfin/dto.rs create mode 100644 backend/src/media_server/jellyfin/mod.rs create mode 100644 backend/src/media_server/mod.rs create mode 100644 backend/src/media_server/playback.rs create mode 100644 backend/src/media_server/playlist_mapper.rs create mode 100644 backend/src/media_server/plex/dto.rs create mode 100644 backend/src/media_server/plex/mod.rs create mode 100644 backend/src/media_server/redaction.rs create mode 100644 backend/src/media_server/test_fixtures.rs create mode 100644 backend/src/media_server/types.rs diff --git a/backend/src/api/endpoints/api_playlist_utils.rs b/backend/src/api/endpoints/api_playlist_utils.rs index c2671ba62..4f47917d4 100644 --- a/backend/src/api/endpoints/api_playlist_utils.rs +++ b/backend/src/api/endpoints/api_playlist_utils.rs @@ -148,6 +148,13 @@ pub(in crate::api::endpoints) async fn get_playlist_for_custom_provider( ) .into_response(); } + InputType::Emby | InputType::Jellyfin | InputType::Plex => { + return ( + axum::http::StatusCode::BAD_REQUEST, + axum::Json(json!({ "error": "Media-server inputs are not supported on this endpoint yet"})), + ) + .into_response(); + } }; if result.is_empty() { let error_strings: Vec = errors.iter().map(ToString::to_string).collect(); diff --git a/backend/src/api/model/metadata_update_manager.rs b/backend/src/api/model/metadata_update_manager.rs index 9bc7bd7e9..e93fb33a1 100644 --- a/backend/src/api/model/metadata_update_manager.rs +++ b/backend/src/api/model/metadata_update_manager.rs @@ -3343,8 +3343,10 @@ impl InputWorker { fn task_needs_provider_connection(task: &UpdateTask, input_type: InputType) -> bool { match task { UpdateTask::ProbeLive { .. } => true, - // Local library probing is fully local and must not depend on provider capacity. - UpdateTask::ProbeStream { .. } => !matches!(input_type, InputType::Library), + // Local library and media-server probing must not depend on IPTV provider capacity. + UpdateTask::ProbeStream { .. } => { + !(matches!(input_type, InputType::Library) || input_type.is_media_server()) + }, // Resolve tasks handle their own probe connection acquisition internally. // This avoids holding a provider connection for the entire duration of // info fetch + TMDB resolve + probe, reducing "provider exhausted" errors. @@ -4257,6 +4259,23 @@ mod tests { assert!(!InputWorker::task_needs_provider_connection(&task, InputType::Library)); } + #[test] + fn task_needs_provider_connection_skips_media_server_probe_stream() { + let task = UpdateTask::ProbeStream { + probe_scope: Arc::from("input_remote"), + unique_id: "u1".to_string(), + url: "http://example.com/stream.mkv".to_string(), + item_type: PlaylistItemType::Video, + reason: ResolveReason::MissingDetails.into(), + delay: 0, + }; + + for input_type in [InputType::Emby, InputType::Jellyfin, InputType::Plex] { + assert!(!InputWorker::task_needs_provider_connection(&task, input_type)); + } + assert!(InputWorker::task_needs_provider_connection(&task, InputType::M3u)); + } + #[test] fn task_needs_provider_connection_keeps_non_library_probe_stream() { let task = UpdateTask::ProbeStream { diff --git a/backend/src/api/model/provider_lineup_manager.rs b/backend/src/api/model/provider_lineup_manager.rs index 80daf121c..a402a26dc 100644 --- a/backend/src/api/model/provider_lineup_manager.rs +++ b/backend/src/api/model/provider_lineup_manager.rs @@ -1039,6 +1039,7 @@ mod tests { input_type: InputType::Xtream, // You can use a default value here max_connections, priority, + media_server: None, aliases: None, headers: HashMap::default(), options: None, diff --git a/backend/src/api/setup_api.rs b/backend/src/api/setup_api.rs index 626a628cb..75b8fe164 100644 --- a/backend/src/api/setup_api.rs +++ b/backend/src/api/setup_api.rs @@ -53,7 +53,7 @@ impl fmt::Debug for SetupWebUserCredentialDto { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { f.debug_struct("SetupWebUserCredentialDto") .field("username", &self.username) - .field("password", &"") + .field("password", &"***") .finish() } } diff --git a/backend/src/media_server/catalog.rs b/backend/src/media_server/catalog.rs new file mode 100644 index 000000000..14a0e2cbe --- /dev/null +++ b/backend/src/media_server/catalog.rs @@ -0,0 +1,307 @@ +use crate::media_server::{ + MediaServerEpisode, MediaServerLibrary, MediaServerLibraryKind, MediaServerCatalogClient, MediaServerError, MediaServerErrorKind, + MediaServerMovie, MediaServerPage, MediaServerPageRequest, +}; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MediaServerCatalogCursor { + pub library_id: String, + pub kind: MediaServerLibraryKind, + pub start: usize, + pub limit: usize, + pub total: Option, + pub fetched: usize, +} + +impl MediaServerCatalogCursor { + pub fn from_page(library: &MediaServerLibrary, page: &MediaServerPage) -> Self { + Self { + library_id: library.reference.library_id.to_string(), + kind: library.kind, + start: page.request.start, + limit: page.request.limit, + total: page.total, + fetched: page.upstream_item_count(), + } + } + + pub fn is_stalled_before_end(&self) -> bool { + self.fetched == 0 && self.total.is_some_and(|total| self.start < total) + } +} + +#[derive(Debug, Copy, Clone, PartialEq, Eq)] +pub struct MediaServerCatalogRefreshPolicy { + pub page_size: usize, +} + +impl Default for MediaServerCatalogRefreshPolicy { + fn default() -> Self { Self { page_size: 100 } } +} + +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct MediaServerCatalogSnapshot { + pub libraries: Vec, + pub movies: Vec, + pub episodes: Vec, + pub unsupported_libraries: Vec, +} + +impl MediaServerCatalogSnapshot { + pub fn item_count(&self) -> usize { self.movies.len() + self.episodes.len() } +} + +#[derive(Debug, Clone, Default)] +pub struct MediaServerCatalogCache { + trusted: Option, +} + +impl MediaServerCatalogCache { + pub fn trusted(&self) -> Option<&MediaServerCatalogSnapshot> { self.trusted.as_ref() } + + pub fn publish(&mut self, snapshot: MediaServerCatalogSnapshot) -> &MediaServerCatalogSnapshot { + self.trusted.insert(snapshot) + } + + pub async fn refresh_or_retain( + &mut self, + client: &C, + policy: MediaServerCatalogRefreshPolicy, + ) -> MediaServerCatalogRefreshOutcome + where + C: MediaServerCatalogClient, + { + match refresh_media_server_catalog_complete_before_publish(client, policy).await { + Ok(snapshot) => { + self.publish(snapshot); + MediaServerCatalogRefreshOutcome::Published + } + Err(error) => MediaServerCatalogRefreshOutcome::Retained { error, retained: self.trusted.is_some() }, + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum MediaServerCatalogRefreshOutcome { + Published, + Retained { error: MediaServerError, retained: bool }, +} + +pub async fn refresh_media_server_catalog_complete_before_publish( + client: &C, + policy: MediaServerCatalogRefreshPolicy, +) -> Result +where + C: MediaServerCatalogClient, +{ + if policy.page_size == 0 { + return Err(MediaServerError::new(MediaServerErrorKind::MediaServerCatalogIncomplete) + .detail("media server catalog page_size must be greater than zero")); + } + + let _server = client.discover().await?; + let libraries = client.list_libraries().await?; + let mut snapshot = MediaServerCatalogSnapshot { libraries: libraries.clone(), ..MediaServerCatalogSnapshot::default() }; + + for library in libraries { + match library.kind { + MediaServerLibraryKind::Movies => { + let mut page_request = MediaServerPageRequest::new(0, policy.page_size); + loop { + let page = client.list_movies(&library.reference, page_request).await?; + validate_page_progress(&library, &page)?; + let next_request = page.next_request(); + snapshot.movies.extend(page.items); + let Some(next) = next_request else { break }; + page_request = next; + } + } + MediaServerLibraryKind::TvShows => { + let mut page_request = MediaServerPageRequest::new(0, policy.page_size); + loop { + let page = client.list_episodes(&library.reference, page_request).await?; + validate_page_progress(&library, &page)?; + let next_request = page.next_request(); + snapshot.episodes.extend(page.items); + let Some(next) = next_request else { break }; + page_request = next; + } + } + MediaServerLibraryKind::Unsupported => snapshot.unsupported_libraries.push(library), + } + } + + Ok(snapshot) +} + +fn validate_page_progress(library: &MediaServerLibrary, page: &MediaServerPage) -> Result<(), MediaServerError> { + let cursor = MediaServerCatalogCursor::from_page(library, page); + if cursor.is_stalled_before_end() { + return Err(MediaServerError::new(MediaServerErrorKind::MediaServerCatalogPageStalled).detail(format!( + "media server catalog page stalled for library kind {:?} at start {}", + cursor.kind, cursor.start + ))); + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::media_server::{ + MediaServerImageRef, MediaServerLibraryRef, MediaServerKind, MediaServerProviderIdHint, + MediaServerResourceResponse, MediaServerStatus, MediaServerStreamRef, MediaServerStreamResponse, + }; + use bytes::Bytes; + use futures::{stream, StreamExt}; + use reqwest::{header::HeaderMap, StatusCode}; + use std::sync::{Arc, Mutex}; + + #[derive(Default)] + struct MockMediaServerCatalogClient { + movie_pages: Mutex, MediaServerError>>>, + episode_pages: Mutex, MediaServerError>>>, + libraries: Vec, + } + + impl MockMediaServerCatalogClient { + fn with_libraries(libraries: Vec) -> Self { + Self { libraries, ..Self::default() } + } + } + + impl MediaServerCatalogClient for MockMediaServerCatalogClient { + async fn discover(&self) -> Result { + Ok(MediaServerStatus { + kind: MediaServerKind::Emby, + server_id: "server-redacted".into(), + display_name: None, + version: None, + owned: None, + }) + } + + async fn list_libraries(&self) -> Result, MediaServerError> { Ok(self.libraries.clone()) } + + async fn list_movies( + &self, + _library: &MediaServerLibraryRef, + _page: MediaServerPageRequest, + ) -> Result, MediaServerError> { + self.movie_pages.lock().expect("lock").remove(0) + } + + async fn list_episodes( + &self, + _library: &MediaServerLibraryRef, + _page: MediaServerPageRequest, + ) -> Result, MediaServerError> { + self.episode_pages.lock().expect("lock").remove(0) + } + + async fn open_stream( + &self, + _stream_ref: &MediaServerStreamRef, + _range: Option<&str>, + ) -> Result { + Ok(empty_stream_response()) + } + + async fn open_image( + &self, + _image_ref: &MediaServerImageRef, + ) -> Result { + Ok(empty_response()) + } + } + + fn empty_response() -> MediaServerResourceResponse { + MediaServerResourceResponse { status: StatusCode::OK, headers: HeaderMap::new(), body: Bytes::new() } + } + + fn empty_stream_response() -> MediaServerStreamResponse { + MediaServerStreamResponse { + status: StatusCode::OK, + headers: HeaderMap::new(), + body: stream::once(async { Ok::(Bytes::new()) }).boxed(), + } + } + + fn movie_library() -> MediaServerLibrary { + MediaServerLibrary { + reference: MediaServerLibraryRef { + input_name: "media_server".into(), + server_id: "server".into(), + library_id: "movies".into(), + }, + name: "Movies".into(), + kind: MediaServerLibraryKind::Movies, + } + } + + fn unsupported_library() -> MediaServerLibrary { + MediaServerLibrary { kind: MediaServerLibraryKind::Unsupported, name: "Music".into(), ..movie_library() } + } + + fn movie(id: &str) -> MediaServerMovie { + MediaServerMovie { + input_name: "media_server".into(), + server_id: "server".into(), + library_id: "movies".into(), + item_id: Arc::::from(id), + title: Arc::::from("Movie Redacted"), + year: None, + source_version_hint: None, + provider_hints: Vec::::new(), + stream_ref: None, + image_ref: None, + } + } + + #[tokio::test] + async fn incomplete_refresh_retains_previous_trusted_snapshot() { + let mut cache = MediaServerCatalogCache::default(); + cache.publish(MediaServerCatalogSnapshot { movies: vec![movie("old")], ..MediaServerCatalogSnapshot::default() }); + + let client = MockMediaServerCatalogClient::with_libraries(vec![movie_library()]); + client.movie_pages.lock().expect("lock").extend([ + Ok(MediaServerPage::new(MediaServerPageRequest::new(0, 1), Some(2), vec![movie("new-1")])), + Err(MediaServerError::new(MediaServerErrorKind::MediaServerUnavailable)), + ]); + + let outcome = cache + .refresh_or_retain(&client, MediaServerCatalogRefreshPolicy { page_size: 1 }) + .await; + + assert!(matches!(outcome, MediaServerCatalogRefreshOutcome::Retained { retained: true, .. })); + assert_eq!(cache.trusted().expect("previous snapshot retained").movies[0].item_id.as_ref(), "old"); + } + + #[tokio::test] + async fn stalled_page_returns_stable_failure() { + let client = MockMediaServerCatalogClient::with_libraries(vec![movie_library()]); + client + .movie_pages + .lock() + .expect("lock") + .push(Ok(MediaServerPage::new(MediaServerPageRequest::new(0, 100), Some(1), vec![]))); + + let error = refresh_media_server_catalog_complete_before_publish(&client, MediaServerCatalogRefreshPolicy::default()) + .await + .expect_err("stalled page should fail"); + + assert_eq!(error.kind, MediaServerErrorKind::MediaServerCatalogPageStalled); + } + + #[tokio::test] + async fn unsupported_library_kind_is_reported_and_not_coerced() { + let client = MockMediaServerCatalogClient::with_libraries(vec![unsupported_library()]); + + let snapshot = refresh_media_server_catalog_complete_before_publish(&client, MediaServerCatalogRefreshPolicy::default()) + .await + .expect("unsupported library should be skipped safely"); + + assert_eq!(snapshot.item_count(), 0); + assert_eq!(snapshot.unsupported_libraries.len(), 1); + } +} diff --git a/backend/src/media_server/client.rs b/backend/src/media_server/client.rs new file mode 100644 index 000000000..7bfab0f14 --- /dev/null +++ b/backend/src/media_server/client.rs @@ -0,0 +1,128 @@ +use crate::media_server::{ + redaction::redact_media_server_text, MediaServerEpisode, MediaServerError, MediaServerErrorKind, MediaServerImageRef, + MediaServerLibrary, MediaServerLibraryRef, MediaServerMovie, MediaServerPage, MediaServerPageRequest, + MediaServerResourceResponse, MediaServerStatus, MediaServerStreamRef, MediaServerStreamResponse, +}; +use reqwest::{ + header::{HeaderMap, HeaderName, HeaderValue}, + Method, RequestBuilder, +}; + +#[allow(async_fn_in_trait)] +pub trait MediaServerCatalogClient: Send + Sync { + async fn discover(&self) -> Result; + + async fn list_libraries(&self) -> Result, MediaServerError>; + + async fn list_movies( + &self, + library: &MediaServerLibraryRef, + page: MediaServerPageRequest, + ) -> Result, MediaServerError>; + + async fn list_episodes( + &self, + library: &MediaServerLibraryRef, + page: MediaServerPageRequest, + ) -> Result, MediaServerError>; + + async fn open_stream( + &self, + stream_ref: &MediaServerStreamRef, + range: Option<&str>, + ) -> Result; + + async fn open_image(&self, image_ref: &MediaServerImageRef) -> Result; +} + +#[derive(Clone)] +pub struct MediaServerHttpClient { + client: reqwest::Client, +} + +impl MediaServerHttpClient { + pub fn new(client: reqwest::Client) -> Self { Self { client } } + + pub fn inner(&self) -> &reqwest::Client { &self.client } + + pub fn request(&self, method: Method, url: &str) -> MediaServerHttpRequestBuilder { + MediaServerHttpRequestBuilder { + safe_url: redact_media_server_text(url), + builder: self.client.request(method, url), + not_found_kind: MediaServerErrorKind::MediaServerItemNotFound, + fallback_kind: MediaServerErrorKind::MediaServerStreamOpenFailed, + } + } +} + +pub struct MediaServerHttpRequestBuilder { + safe_url: String, + builder: RequestBuilder, + not_found_kind: MediaServerErrorKind, + fallback_kind: MediaServerErrorKind, +} + +impl MediaServerHttpRequestBuilder { + pub fn safe_url(&self) -> &str { &self.safe_url } + + pub fn header(mut self, key: HeaderName, value: HeaderValue) -> Self { + self.builder = self.builder.header(key, value); + self + } + + pub fn headers(mut self, headers: HeaderMap) -> Self { + self.builder = self.builder.headers(headers); + self + } + + pub fn error_kinds(mut self, not_found_kind: MediaServerErrorKind, fallback_kind: MediaServerErrorKind) -> Self { + self.not_found_kind = not_found_kind; + self.fallback_kind = fallback_kind; + self + } + + pub fn discovery_errors(self) -> Self { + self.error_kinds(MediaServerErrorKind::MediaServerUnavailable, MediaServerErrorKind::MediaServerDiscoveryFailed) + } + + pub fn catalog_errors(self) -> Self { + self.error_kinds( + MediaServerErrorKind::MediaServerLibraryUnavailable, + MediaServerErrorKind::MediaServerCatalogDecodeFailed, + ) + } + + pub fn playback_errors(self) -> Self { + self.error_kinds(MediaServerErrorKind::MediaServerItemNotFound, MediaServerErrorKind::MediaServerStreamOpenFailed) + } + + pub async fn send(self) -> Result { + self.builder.send().await.map_err(|err| { + MediaServerError::from_reqwest_error_with_fallback(&err, self.not_found_kind, self.fallback_kind) + .detail(format!("request {} failed", self.safe_url)) + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn media_server_http_request_builder_keeps_safe_url_redacted() { + let client = MediaServerHttpClient::new(reqwest::Client::new()); + let request = client.request(Method::GET, "https://media.example.invalid/video?api_key=secret"); + + assert!(!request.safe_url().contains("secret")); + assert!(request.safe_url().contains("api_key=***")); + } + + #[test] + fn media_server_http_request_builder_can_select_catalog_error_context() { + let client = MediaServerHttpClient::new(reqwest::Client::new()); + let request = client.request(Method::GET, "https://media.example.invalid/libraries").catalog_errors(); + + assert_eq!(request.not_found_kind, MediaServerErrorKind::MediaServerLibraryUnavailable); + assert_eq!(request.fallback_kind, MediaServerErrorKind::MediaServerCatalogDecodeFailed); + } +} diff --git a/backend/src/media_server/emby/dto.rs b/backend/src/media_server/emby/dto.rs new file mode 100644 index 000000000..678a0bd3e --- /dev/null +++ b/backend/src/media_server/emby/dto.rs @@ -0,0 +1,88 @@ +use serde::Deserialize; +use std::collections::HashMap; + +#[derive(Debug, Clone, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "PascalCase")] +pub struct EmbyPublicSystemInfoDto { + pub id: Option, + pub server_name: Option, + pub version: Option, +} + +#[derive(Debug, Clone, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "PascalCase", bound(deserialize = "T: Deserialize<'de>"))] +pub struct EmbyItemsPageDto { + #[serde(default)] + pub items: Vec, + #[serde(default)] + pub total_record_count: usize, + #[serde(default)] + pub start_index: usize, +} + +#[derive(Debug, Clone, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "PascalCase")] +pub struct EmbyViewDto { + pub id: String, + pub name: Option, + pub collection_type: Option, + pub type_: Option, +} + +#[derive(Debug, Clone, Deserialize, PartialEq)] +#[serde(rename_all = "PascalCase")] +pub struct EmbyItemDto { + pub id: String, + pub name: Option, + pub type_: Option, + pub production_year: Option, + pub parent_id: Option, + pub series_id: Option, + pub series_name: Option, + pub parent_index_number: Option, + pub index_number: Option, + #[serde(default)] + pub provider_ids: HashMap, + #[serde(default)] + pub image_tags: HashMap, + #[serde(default)] + pub media_sources: Vec, + // Parsed only so the boundary can explicitly ignore it by default. + pub path: Option, + pub user_data: Option, +} + +#[derive(Debug, Clone, Deserialize, PartialEq)] +#[serde(rename_all = "PascalCase")] +pub struct EmbyMediaSourceDto { + pub id: Option, + pub container: Option, + pub path: Option, + pub supports_direct_play: Option, + pub supports_direct_stream: Option, +} + +#[derive(Debug, Clone, Deserialize, PartialEq)] +#[serde(rename_all = "PascalCase")] +pub struct EmbyPlaybackInfoDto { + #[serde(default)] + pub media_sources: Vec, + pub play_session_id: Option, +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::media_server::test_fixtures::EMBY_ITEMS_PAGE_JSON; + + #[test] + fn parses_emby_item_page_and_keeps_edge_only_fields_at_boundary() { + let page: EmbyItemsPageDto = serde_json::from_str(EMBY_ITEMS_PAGE_JSON).expect("fixture parses"); + + assert_eq!(page.total_record_count, 1); + assert_eq!(page.items[0].id, "item-redacted-1"); + assert_eq!(page.items[0].provider_ids.get("Tmdb").map(String::as_str), Some("12345")); + assert!(page.items[0].path.as_deref().is_some_and(|path| path.contains("/redacted/"))); + assert!(page.items[0].user_data.is_some()); + } +} diff --git a/backend/src/media_server/emby/mod.rs b/backend/src/media_server/emby/mod.rs new file mode 100644 index 000000000..a07dce5c0 --- /dev/null +++ b/backend/src/media_server/emby/mod.rs @@ -0,0 +1 @@ +pub mod dto; diff --git a/backend/src/media_server/errors.rs b/backend/src/media_server/errors.rs new file mode 100644 index 000000000..d4f7c282e --- /dev/null +++ b/backend/src/media_server/errors.rs @@ -0,0 +1,190 @@ +use crate::media_server::redaction::redact_media_server_text; +use reqwest::StatusCode; +use std::{error::Error, fmt}; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum MediaServerErrorKind { + MediaServerAuthDenied, + MediaServerUnavailable, + MediaServerLibraryUnavailable, + MediaServerLibraryTypeUnsupported, + MediaServerCatalogDecodeFailed, + MediaServerCatalogPageStalled, + MediaServerCatalogIncomplete, + MediaServerItemNotFound, + NoDirectPlayableMediaServerSource, + MediaServerStreamOpenFailed, + MediaServerRateLimited, + MediaServerDiscoveryFailed, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MediaServerError { + pub kind: MediaServerErrorKind, + pub provider: Option<&'static str>, + pub status: Option, + detail: Option, +} + +impl MediaServerError { + pub fn new(kind: MediaServerErrorKind) -> Self { + Self { + kind, + provider: None, + status: None, + detail: None, + } + } + + pub fn provider(mut self, provider: &'static str) -> Self { + self.provider = Some(provider); + self + } + + pub fn status(mut self, status: StatusCode) -> Self { + self.status = Some(status); + self + } + + pub fn detail(mut self, detail: impl AsRef) -> Self { + let redacted = redact_media_server_text(detail.as_ref()); + if !redacted.trim().is_empty() { + self.detail = Some(redacted); + } + self + } + + pub fn detail_text(&self) -> Option<&str> { self.detail.as_deref() } + + pub fn from_http_status(status: StatusCode) -> Self { + Self::from_http_status_with_fallback( + status, + MediaServerErrorKind::MediaServerItemNotFound, + MediaServerErrorKind::MediaServerStreamOpenFailed, + ) + } + + pub fn from_http_status_with_fallback( + status: StatusCode, + not_found_kind: MediaServerErrorKind, + fallback_kind: MediaServerErrorKind, + ) -> Self { + let kind = if matches!(status, StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN) { + MediaServerErrorKind::MediaServerAuthDenied + } else if status == StatusCode::TOO_MANY_REQUESTS { + MediaServerErrorKind::MediaServerRateLimited + } else if status == StatusCode::NOT_FOUND { + not_found_kind + } else { + fallback_kind + }; + Self::new(kind).status(status) + } + + pub fn from_reqwest_error(err: &reqwest::Error) -> Self { + Self::from_reqwest_error_with_fallback( + err, + MediaServerErrorKind::MediaServerItemNotFound, + MediaServerErrorKind::MediaServerStreamOpenFailed, + ) + } + + pub fn from_reqwest_error_with_fallback( + err: &reqwest::Error, + not_found_kind: MediaServerErrorKind, + fallback_kind: MediaServerErrorKind, + ) -> Self { + let kind = if err.is_timeout() || err.is_connect() { + MediaServerErrorKind::MediaServerUnavailable + } else if err.is_decode() { + fallback_kind + } else if let Some(status) = err.status() { + return Self::from_http_status_with_fallback(status, not_found_kind, fallback_kind).detail(err.to_string()); + } else { + fallback_kind + }; + Self::new(kind).detail(err.to_string()) + } +} + +impl fmt::Display for MediaServerError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "{:?}", self.kind)?; + if let Some(provider) = self.provider { + write!(f, " provider={provider}")?; + } + if let Some(status) = self.status { + write!(f, " status={status}")?; + } + if let Some(detail) = self.detail.as_deref() { + write!(f, " detail={detail}")?; + } + Ok(()) + } +} + +impl Error for MediaServerError {} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn http_status_maps_to_stable_failure_kinds() { + assert_eq!( + MediaServerError::from_http_status(StatusCode::UNAUTHORIZED).kind, + MediaServerErrorKind::MediaServerAuthDenied + ); + assert_eq!( + MediaServerError::from_http_status(StatusCode::TOO_MANY_REQUESTS).kind, + MediaServerErrorKind::MediaServerRateLimited + ); + assert_eq!( + MediaServerError::from_http_status(StatusCode::NOT_FOUND).kind, + MediaServerErrorKind::MediaServerItemNotFound + ); + } + + #[test] + fn http_status_mapping_accepts_operation_specific_fallbacks() { + assert_eq!( + MediaServerError::from_http_status_with_fallback( + StatusCode::NOT_FOUND, + MediaServerErrorKind::MediaServerLibraryUnavailable, + MediaServerErrorKind::MediaServerCatalogDecodeFailed, + ) + .kind, + MediaServerErrorKind::MediaServerLibraryUnavailable + ); + assert_eq!( + MediaServerError::from_http_status_with_fallback( + StatusCode::INTERNAL_SERVER_ERROR, + MediaServerErrorKind::MediaServerLibraryUnavailable, + MediaServerErrorKind::MediaServerCatalogDecodeFailed, + ) + .kind, + MediaServerErrorKind::MediaServerCatalogDecodeFailed + ); + assert_eq!( + MediaServerError::from_http_status_with_fallback( + StatusCode::FORBIDDEN, + MediaServerErrorKind::MediaServerLibraryUnavailable, + MediaServerErrorKind::MediaServerCatalogDecodeFailed, + ) + .kind, + MediaServerErrorKind::MediaServerAuthDenied + ); + } + + #[test] + fn detail_is_redacted_before_display() { + let err = MediaServerError::new(MediaServerErrorKind::MediaServerStreamOpenFailed) + .provider("plex") + .detail("https://media.example.invalid/stream?X-Plex-Token=secret-token&api_key=secret-key"); + let rendered = err.to_string(); + + assert!(!rendered.contains("secret-token")); + assert!(!rendered.contains("secret-key")); + assert!(rendered.contains("***")); + } +} diff --git a/backend/src/media_server/jellyfin/dto.rs b/backend/src/media_server/jellyfin/dto.rs new file mode 100644 index 000000000..708be2364 --- /dev/null +++ b/backend/src/media_server/jellyfin/dto.rs @@ -0,0 +1,26 @@ +//! Jellyfin uses Emby-compatible JSON shapes for the MVP endpoints Tuliprox needs. +//! +//! Keep aliases in this module so caller code still depends on a provider-specific +//! seam and can diverge safely when Jellyfin fields differ. + +pub type JellyfinPublicSystemInfoDto = crate::media_server::emby::dto::EmbyPublicSystemInfoDto; +pub type JellyfinItemsPageDto = crate::media_server::emby::dto::EmbyItemsPageDto; +pub type JellyfinViewDto = crate::media_server::emby::dto::EmbyViewDto; +pub type JellyfinItemDto = crate::media_server::emby::dto::EmbyItemDto; +pub type JellyfinMediaSourceDto = crate::media_server::emby::dto::EmbyMediaSourceDto; +pub type JellyfinPlaybackInfoDto = crate::media_server::emby::dto::EmbyPlaybackInfoDto; + +#[cfg(test)] +mod tests { + use super::*; + use crate::media_server::test_fixtures::JELLYFIN_VIEWS_JSON; + + #[test] + fn parses_jellyfin_views_through_provider_specific_aliases() { + let page: JellyfinItemsPageDto = serde_json::from_str(JELLYFIN_VIEWS_JSON).expect("fixture parses"); + + assert_eq!(page.items.len(), 2); + assert_eq!(page.items[0].collection_type.as_deref(), Some("movies")); + assert_eq!(page.items[1].collection_type.as_deref(), Some("tvshows")); + } +} diff --git a/backend/src/media_server/jellyfin/mod.rs b/backend/src/media_server/jellyfin/mod.rs new file mode 100644 index 000000000..a07dce5c0 --- /dev/null +++ b/backend/src/media_server/jellyfin/mod.rs @@ -0,0 +1 @@ +pub mod dto; diff --git a/backend/src/media_server/mod.rs b/backend/src/media_server/mod.rs new file mode 100644 index 000000000..969f07cba --- /dev/null +++ b/backend/src/media_server/mod.rs @@ -0,0 +1,27 @@ +//! Media-server anti-corruption layer. +//! +//! This module keeps Emby/Jellyfin/Plex wire DTOs and transport concerns at the +//! Source Acquisition boundary. Playlist Curation and Stream Brokerage should +//! depend on the typed concepts exported here instead of provider-specific DTOs. + +pub mod catalog; +pub mod client; +pub mod emby; +pub mod errors; +pub mod jellyfin; +pub mod playback; +pub mod playlist_mapper; +pub mod plex; +pub mod redaction; +pub mod types; + +#[cfg(test)] +pub mod test_fixtures; + +pub use catalog::*; +pub use client::*; +pub use errors::*; +pub use playback::*; +pub use playlist_mapper::*; +pub use redaction::*; +pub use types::*; diff --git a/backend/src/media_server/playback.rs b/backend/src/media_server/playback.rs new file mode 100644 index 000000000..22b993780 --- /dev/null +++ b/backend/src/media_server/playback.rs @@ -0,0 +1,335 @@ +use crate::media_server::{ + BoxedMediaServerStream, MediaServerCatalogClient, MediaServerError, MediaServerErrorKind, MediaServerImageRef, + MediaServerResourceResponse, MediaServerStreamRef, MediaServerStreamResponse, +}; +use bytes::Bytes; +use futures::{stream, StreamExt}; +use reqwest::{ + header::{ + HeaderMap, ACCEPT_RANGES, CONTENT_LENGTH, CONTENT_RANGE, CONTENT_TYPE, ETAG, LAST_MODIFIED, + }, + StatusCode, +}; +use shared::model::{InputType, PlaylistItemType}; +use std::{fmt, sync::Arc}; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum PlaybackOrigin { + Provider, + LocalLibrary, + MediaServer(MediaServerStreamRef), +} + +pub struct MediaServerProxyResponse { + pub status: StatusCode, + pub headers: HeaderMap, + pub body: BoxedMediaServerStream, +} + +impl fmt::Debug for MediaServerProxyResponse { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("MediaServerProxyResponse") + .field("status", &self.status) + .field("headers", &self.headers) + .field("body", &"") + .finish() + } +} + +pub fn classify_playback_origin( + input_type: InputType, + item_type: PlaylistItemType, + input_name: &Arc, + item_url: &str, +) -> Result { + if input_type.is_media_server() || item_url.starts_with("media-server://") { + return parse_media_server_stream_ref(input_name, item_url).map(PlaybackOrigin::MediaServer); + } + + if matches!(item_type, PlaylistItemType::LocalVideo | PlaylistItemType::LocalSeries) { + return Ok(PlaybackOrigin::LocalLibrary); + } + + Ok(PlaybackOrigin::Provider) +} + +pub async fn media_server_stream_response( + client: &C, + stream_ref: &MediaServerStreamRef, + range: Option<&str>, +) -> Result +where + C: MediaServerCatalogClient, +{ + let response = client.open_stream(stream_ref, range).await?; + Ok(media_server_stream_to_proxy_response(response)) +} + +pub async fn media_server_image_response( + client: &C, + image_ref: &MediaServerImageRef, +) -> Result +where + C: MediaServerCatalogClient, +{ + let response = client.open_image(image_ref).await?; + Ok(media_server_resource_to_proxy_response(response)) +} + +fn media_server_resource_to_proxy_response(response: MediaServerResourceResponse) -> MediaServerProxyResponse { + MediaServerProxyResponse { + status: response.status, + headers: safe_media_server_response_headers(&response.headers), + body: single_chunk_stream(response.body), + } +} + +fn media_server_stream_to_proxy_response(response: MediaServerStreamResponse) -> MediaServerProxyResponse { + MediaServerProxyResponse { + status: response.status, + headers: safe_media_server_response_headers(&response.headers), + body: response.body, + } +} + +fn single_chunk_stream(body: Bytes) -> BoxedMediaServerStream { + stream::once(async move { Ok::(body) }).boxed() +} + +pub fn safe_media_server_response_headers(headers: &HeaderMap) -> HeaderMap { + let mut safe = HeaderMap::new(); + for name in [CONTENT_TYPE, CONTENT_LENGTH, CONTENT_RANGE, ACCEPT_RANGES, ETAG, LAST_MODIFIED] { + if let Some(value) = headers.get(&name) { + safe.insert(name, value.clone()); + } + } + safe +} + +pub fn parse_media_server_stream_ref(input_name: &Arc, item_url: &str) -> Result { + let Some(rest) = item_url.strip_prefix("media-server://") else { + return Err(MediaServerError::new(MediaServerErrorKind::MediaServerStreamOpenFailed) + .detail("playlist item is not a media server URL")); + }; + let (path, query) = rest.split_once('?').unwrap_or((rest, "")); + let parts: Vec = path.split('/').map(unescape_internal_url_component).collect(); + if parts.len() < 3 { + return Err(MediaServerError::new(MediaServerErrorKind::MediaServerStreamOpenFailed) + .detail("media server URL is missing required path parts")); + } + + match parts[0].as_str() { + "unavailable" => Err( + MediaServerError::new(MediaServerErrorKind::NoDirectPlayableMediaServerSource) + .detail("media server item has no direct playable source"), + ), + "emby" => Ok(MediaServerStreamRef::Emby { + input_name: input_name.clone(), + server_id: Arc::::from(parts[1].as_str()), + item_id: Arc::::from(parts[2].as_str()), + media_source_id: query_value(query, "media_source_id").map(Arc::::from), + }), + "jellyfin" => Ok(MediaServerStreamRef::Jellyfin { + input_name: input_name.clone(), + server_id: Arc::::from(parts[1].as_str()), + item_id: Arc::::from(parts[2].as_str()), + media_source_id: query_value(query, "media_source_id").map(Arc::::from), + }), + "plex" => Ok(MediaServerStreamRef::Plex { + input_name: input_name.clone(), + server_id: Arc::::from(parts[1].as_str()), + rating_key: Arc::::from(parts[2].as_str()), + part_key: query_value(query, "part_key") + .map(Arc::::from) + .ok_or_else(|| { + MediaServerError::new(MediaServerErrorKind::NoDirectPlayableMediaServerSource) + .detail("plex media-server URL is missing part_key") + })?, + }), + _ => Err(MediaServerError::new(MediaServerErrorKind::MediaServerStreamOpenFailed) + .detail("unsupported media server URL scheme")), + } +} + +fn query_value(query: &str, key: &str) -> Option { + url::form_urlencoded::parse(query.as_bytes()) + .find_map(|(name, value)| (name == key).then(|| value.into_owned())) +} + +fn unescape_internal_url_component(value: &str) -> String { + let bytes = value.as_bytes(); + let mut decoded = Vec::with_capacity(bytes.len()); + let mut i = 0; + while i < bytes.len() { + if bytes[i] == b'%' && i + 2 < bytes.len() { + if let Some(byte) = decode_hex_byte(bytes[i + 1], bytes[i + 2]) { + decoded.push(byte); + i += 3; + continue; + } + } + decoded.push(bytes[i]); + i += 1; + } + String::from_utf8_lossy(&decoded).into_owned() +} + +fn decode_hex_byte(high: u8, low: u8) -> Option { + Some(hex_value(high)? << 4 | hex_value(low)?) +} + +fn hex_value(value: u8) -> Option { + match value { + b'0'..=b'9' => Some(value - b'0'), + b'a'..=b'f' => Some(value - b'a' + 10), + b'A'..=b'F' => Some(value - b'A' + 10), + _ => None, + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::media_server::{ + MediaServerEpisode, MediaServerLibrary, MediaServerLibraryRef, MediaServerMovie, MediaServerPage, MediaServerPageRequest, + MediaServerStatus, + }; + use futures::{stream, StreamExt}; + use reqwest::header::{HeaderValue, AUTHORIZATION}; + use std::sync::{Mutex, Arc as StdArc}; + + #[derive(Default)] + struct MockPlaybackClient { + seen_range: Mutex>, + stream_error: Option, + } + + impl MediaServerCatalogClient for MockPlaybackClient { + async fn discover(&self) -> Result { unreachable!() } + async fn list_libraries(&self) -> Result, MediaServerError> { unreachable!() } + async fn list_movies( + &self, + _library: &MediaServerLibraryRef, + _page: MediaServerPageRequest, + ) -> Result, MediaServerError> { + unreachable!() + } + async fn list_episodes( + &self, + _library: &MediaServerLibraryRef, + _page: MediaServerPageRequest, + ) -> Result, MediaServerError> { + unreachable!() + } + + async fn open_stream( + &self, + _stream_ref: &MediaServerStreamRef, + range: Option<&str>, + ) -> Result { + *self.seen_range.lock().expect("lock") = range.map(ToOwned::to_owned); + if let Some(error) = self.stream_error.clone() { + return Err(error); + } + let mut headers = HeaderMap::new(); + headers.insert(CONTENT_TYPE, HeaderValue::from_static("video/mp4")); + headers.insert(CONTENT_RANGE, HeaderValue::from_static("bytes 0-1023/2048")); + headers.insert(CONTENT_LENGTH, HeaderValue::from_static("1024")); + headers.insert(ACCEPT_RANGES, HeaderValue::from_static("bytes")); + headers.insert(AUTHORIZATION, HeaderValue::from_static("Bearer should-not-leak")); + Ok(MediaServerStreamResponse { + status: StatusCode::PARTIAL_CONTENT, + headers, + body: stream::once(async { Ok::(Bytes::from_static(b"data")) }).boxed(), + }) + } + + async fn open_image(&self, _image_ref: &MediaServerImageRef) -> Result { + Ok(MediaServerResourceResponse { status: StatusCode::OK, headers: HeaderMap::new(), body: Bytes::new() }) + } + } + + #[tokio::test] + async fn media_server_stream_response_forwards_range_and_filters_headers() { + let client = MockPlaybackClient::default(); + let stream_ref = MediaServerStreamRef::Plex { + input_name: "media_server".into(), + server_id: "server".into(), + rating_key: "rating".into(), + part_key: "/library/parts/redacted/file.mkv".into(), + }; + + let response = media_server_stream_response(&client, &stream_ref, Some("bytes=0-1023")) + .await + .expect("stream opens"); + + assert_eq!(client.seen_range.lock().expect("lock").as_deref(), Some("bytes=0-1023")); + assert_eq!(response.status, StatusCode::PARTIAL_CONTENT); + assert_eq!(response.headers.get(CONTENT_RANGE).and_then(|v| v.to_str().ok()), Some("bytes 0-1023/2048")); + assert!(response.headers.get(AUTHORIZATION).is_none()); + + let chunks = response.body.collect::>().await; + assert_eq!(chunks.len(), 1); + assert_eq!(chunks[0].as_ref().map(Bytes::as_ref), Ok(b"data".as_slice())); + } + + #[test] + fn classifies_media_server_local_and_provider_origins() { + let input_name = StdArc::::from("media_server"); + let media_server = classify_playback_origin( + InputType::Plex, + PlaylistItemType::Video, + &input_name, + "media-server://plex/server/rating?part_key=%2Flibrary%2Fparts%2Fredacted%2Ffile.mkv", + ) + .expect("media_server ref parses"); + assert!(matches!(media_server, PlaybackOrigin::MediaServer(MediaServerStreamRef::Plex { .. }))); + + let local = classify_playback_origin(InputType::Library, PlaylistItemType::LocalVideo, &input_name, "file:///tmp/a.mkv") + .expect("local classifies"); + assert_eq!(local, PlaybackOrigin::LocalLibrary); + + let encoded = parse_media_server_stream_ref( + &input_name, + "media-server://emby/server%2Fone/item%3Fone?media_source_id=media%2Fsource%3Fone", + ) + .expect("encoded media_server ref parses"); + assert_eq!( + encoded, + MediaServerStreamRef::Emby { + input_name: input_name.clone(), + server_id: "server/one".into(), + item_id: "item?one".into(), + media_source_id: Some("media/source?one".into()), + } + ); + + let unavailable = parse_media_server_stream_ref(&input_name, "media-server://unavailable/server/library/item") + .expect_err("unavailable sentinel should map to a stable playback error"); + assert_eq!(unavailable.kind, MediaServerErrorKind::NoDirectPlayableMediaServerSource); + + let provider = classify_playback_origin(InputType::M3u, PlaylistItemType::Live, &input_name, "http://example.invalid/live") + .expect("provider classifies"); + assert_eq!(provider, PlaybackOrigin::Provider); + } + + #[tokio::test] + async fn media_server_auth_denied_stays_media_server_error() { + let client = MockPlaybackClient { + stream_error: Some(MediaServerError::new(MediaServerErrorKind::MediaServerAuthDenied)), + ..MockPlaybackClient::default() + }; + let stream_ref = MediaServerStreamRef::Emby { + input_name: "media_server".into(), + server_id: "server".into(), + item_id: "item".into(), + media_source_id: None, + }; + + let error = media_server_stream_response(&client, &stream_ref, None) + .await + .expect_err("auth denied should fail"); + + assert_eq!(error.kind, MediaServerErrorKind::MediaServerAuthDenied); + } +} diff --git a/backend/src/media_server/playlist_mapper.rs b/backend/src/media_server/playlist_mapper.rs new file mode 100644 index 000000000..57cda45f4 --- /dev/null +++ b/backend/src/media_server/playlist_mapper.rs @@ -0,0 +1,249 @@ +use crate::media_server::{MediaServerCatalogSnapshot, MediaServerEpisode, MediaServerMovie, MediaServerStreamRef}; +use shared::{ + model::{ + EpisodeStreamProperties, PlaylistGroup, PlaylistItem, PlaylistItemHeader, PlaylistItemType, StreamProperties, + VideoStreamProperties, XtreamCluster, + }, + utils::{generate_provider_playlist_uuid, Internable}, +}; +use std::{fmt::Write as _, sync::Arc}; + +pub fn media_server_catalog_snapshot_to_playlist(snapshot: &MediaServerCatalogSnapshot) -> Vec { + let mut groups = Vec::new(); + + if !snapshot.movies.is_empty() { + groups.push(PlaylistGroup { + id: 1, + title: "Media Server Movies".intern(), + channels: snapshot.movies.iter().map(media_server_movie_to_playlist_item).collect(), + xtream_cluster: XtreamCluster::Video, + }); + } + + if !snapshot.episodes.is_empty() { + groups.push(PlaylistGroup { + id: next_group_id(groups.len()), + title: "Media Server Series".intern(), + channels: snapshot.episodes.iter().map(media_server_episode_to_playlist_item).collect(), + xtream_cluster: XtreamCluster::Series, + }); + } + + groups +} + +fn next_group_id(group_count: usize) -> u32 { u32::try_from(group_count.saturating_add(1)).unwrap_or(u32::MAX) } + +fn media_server_movie_to_playlist_item(movie: &MediaServerMovie) -> PlaylistItem { + let stable_id = stable_media_server_item_id(&movie.server_id, &movie.library_id, &movie.item_id, "movie"); + let url = movie.stream_ref.as_ref().map_or_else( + || { + format!( + "media-server://unavailable/{}/{}/{}", + escape_internal_url_component(&movie.server_id), + escape_internal_url_component(&movie.library_id), + escape_internal_url_component(&movie.item_id) + ) + }, + media_server_stream_ref_to_internal_url, + ); + let uuid = generate_provider_playlist_uuid(&movie.input_name, &stable_id, PlaylistItemType::Video); + + PlaylistItem { + header: PlaylistItemHeader { + uuid, + id: stable_id.intern(), + name: movie.title.clone(), + title: movie.title.clone(), + group: "Media Server Movies".intern(), + url: url.intern(), + input_name: movie.input_name.clone(), + xtream_cluster: XtreamCluster::Video, + item_type: PlaylistItemType::Video, + additional_properties: Some(StreamProperties::Video(Box::new(VideoStreamProperties { + name: movie.title.clone(), + stream_id: 0, + stream_icon: "".intern(), + direct_source: "".intern(), + category_id: 0, + custom_sid: None, + added: movie.source_version_hint.clone().unwrap_or_else(|| "".intern()), + container_extension: "".intern(), + rating: None, + rating_5based: None, + stream_type: Some("movie".intern()), + trailer: None, + tmdb: None, + is_adult: 0, + details: None, + }))), + ..PlaylistItemHeader::default() + }, + } +} + +fn media_server_episode_to_playlist_item(episode: &MediaServerEpisode) -> PlaylistItem { + let stable_id = stable_media_server_item_id(&episode.server_id, &episode.library_id, &episode.item_id, "episode"); + let url = episode.stream_ref.as_ref().map_or_else( + || { + format!( + "media-server://unavailable/{}/{}/{}", + escape_internal_url_component(&episode.server_id), + escape_internal_url_component(&episode.library_id), + escape_internal_url_component(&episode.item_id) + ) + }, + media_server_stream_ref_to_internal_url, + ); + let uuid = generate_provider_playlist_uuid(&episode.input_name, &stable_id, PlaylistItemType::Series); + let title = if episode.title.is_empty() { + episode.series_title.clone().unwrap_or_else(|| "Media Server Episode".intern()) + } else { + episode.title.clone() + }; + + PlaylistItem { + header: PlaylistItemHeader { + uuid, + id: stable_id.intern(), + name: title.clone(), + title, + group: "Media Server Series".intern(), + parent_code: episode.series_id.clone().unwrap_or_else(|| "".intern()), + url: url.intern(), + input_name: episode.input_name.clone(), + xtream_cluster: XtreamCluster::Series, + item_type: PlaylistItemType::Series, + additional_properties: Some(StreamProperties::Episode(Box::new(EpisodeStreamProperties { + episode_id: 0, + episode: episode.episode.unwrap_or_default(), + season: episode.season.unwrap_or_default(), + added: episode.source_version_hint.clone(), + release_date: None, + series_release_date: None, + tmdb: None, + movie_image: "".intern(), + container_extension: "".intern(), + video: None, + audio: None, + }))), + ..PlaylistItemHeader::default() + }, + } +} + +fn stable_media_server_item_id(server_id: &Arc, library_id: &Arc, item_id: &Arc, kind: &str) -> String { + format!("media-server:{server_id}:{library_id}:{kind}:{item_id}") +} + +pub fn media_server_stream_ref_to_internal_url(stream_ref: &MediaServerStreamRef) -> String { + match stream_ref { + MediaServerStreamRef::Emby { server_id, item_id, media_source_id, .. } => { + format!( + "media-server://emby/{}/{}{}", + escape_internal_url_component(server_id), + escape_internal_url_component(item_id), + media_source_id + .as_ref() + .map(|id| format!("?media_source_id={}", escape_internal_url_component(id))) + .unwrap_or_default() + ) + } + MediaServerStreamRef::Jellyfin { server_id, item_id, media_source_id, .. } => { + format!( + "media-server://jellyfin/{}/{}{}", + escape_internal_url_component(server_id), + escape_internal_url_component(item_id), + media_source_id + .as_ref() + .map(|id| format!("?media_source_id={}", escape_internal_url_component(id))) + .unwrap_or_default() + ) + } + MediaServerStreamRef::Plex { server_id, rating_key, part_key, .. } => format!( + "media-server://plex/{}/{}?part_key={}", + escape_internal_url_component(server_id), + escape_internal_url_component(rating_key), + escape_internal_url_component(part_key) + ), + } +} + +fn escape_internal_url_component(value: &str) -> String { + let mut encoded = String::with_capacity(value.len()); + for byte in value.bytes() { + if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'.' | b'_' | b'~') { + encoded.push(byte as char); + } else { + let _ = write!(encoded, "%{byte:02X}"); + } + } + encoded +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::media_server::{MediaServerCatalogSnapshot, MediaServerProviderIdHint}; + + fn movie() -> MediaServerMovie { + MediaServerMovie { + input_name: "media_server".into(), + server_id: "server/one".into(), + library_id: "movies".into(), + item_id: "item?one plus+space".into(), + title: "Movie".into(), + year: Some(2024), + source_version_hint: None, + provider_hints: Vec::::new(), + stream_ref: Some(MediaServerStreamRef::Emby { + input_name: "media_server".into(), + server_id: "server/one".into(), + item_id: "item?one plus+space".into(), + media_source_id: Some("media/source".into()), + }), + image_ref: None, + } + } + + fn episode() -> MediaServerEpisode { + MediaServerEpisode { + input_name: "media_server".into(), + server_id: "server".into(), + library_id: "shows".into(), + item_id: "episode".into(), + series_id: Some("series".into()), + series_title: Some("Show".into()), + title: "Episode".into(), + season: Some(1), + episode: Some(2), + source_version_hint: None, + provider_hints: Vec::::new(), + stream_ref: Some(MediaServerStreamRef::Plex { + input_name: "media_server".into(), + server_id: "server".into(), + rating_key: "rating".into(), + part_key: "/library/parts/redacted/file.mkv".into(), + }), + image_ref: None, + } + } + + #[test] + fn maps_media_server_movies_and_episodes_to_playlist_groups_without_virtual_ids() { + let groups = media_server_catalog_snapshot_to_playlist(&MediaServerCatalogSnapshot { + movies: vec![movie()], + episodes: vec![episode()], + ..MediaServerCatalogSnapshot::default() + }); + + assert_eq!(groups.len(), 2); + assert_eq!(groups[0].xtream_cluster, XtreamCluster::Video); + assert_eq!(groups[0].channels[0].header.item_type, PlaylistItemType::Video); + assert_eq!(groups[0].channels[0].header.virtual_id, 0); + assert!(groups[0].channels[0].header.id.starts_with("media-server:server/one:movies:movie:item?one plus+space")); + assert!(groups[0].channels[0].header.url.contains("media-server://emby/server%2Fone/item%3Fone%20plus%2Bspace")); + assert_eq!(groups[1].channels[0].header.item_type, PlaylistItemType::Series); + assert!(groups[1].channels[0].header.url.contains("part_key=%2Flibrary%2Fparts%2Fredacted%2Ffile.mkv")); + } +} diff --git a/backend/src/media_server/plex/dto.rs b/backend/src/media_server/plex/dto.rs new file mode 100644 index 000000000..9c07eca9e --- /dev/null +++ b/backend/src/media_server/plex/dto.rs @@ -0,0 +1,198 @@ +use serde::Deserialize; + +#[derive(Debug, Clone, Deserialize, PartialEq)] +#[serde(rename = "MediaContainer")] +pub struct PlexResourcesDto { + #[serde(rename = "Device", default)] + pub devices: Vec, +} + +#[derive(Debug, Clone, Deserialize, PartialEq)] +pub struct PlexResourceDto { + #[serde(rename = "@name")] + pub name: Option, + #[serde(rename = "@product")] + pub product: Option, + #[serde(rename = "@productVersion")] + pub product_version: Option, + #[serde(rename = "@clientIdentifier")] + pub client_identifier: Option, + #[serde(rename = "@machineIdentifier")] + pub machine_identifier: Option, + #[serde(rename = "@owned", default)] + pub owned: Option, + #[serde(rename = "@accessToken")] + pub access_token: Option, + #[serde(rename = "Connection", default)] + pub connections: Vec, +} + +#[derive(Debug, Clone, Deserialize, PartialEq)] +pub struct PlexConnectionDto { + #[serde(rename = "@protocol")] + pub protocol: Option, + #[serde(rename = "@uri")] + pub uri: Option, + #[serde(rename = "@local", default)] + pub local: Option, + #[serde(rename = "@relay", default)] + pub relay: Option, +} + +#[derive(Debug, Clone, Deserialize, PartialEq)] +#[serde(rename = "MediaContainer")] +pub struct PlexSectionsDto { + #[serde(rename = "Directory", default)] + pub directories: Vec, +} + +#[derive(Debug, Clone, Deserialize, PartialEq)] +pub struct PlexSectionDto { + #[serde(rename = "@key")] + pub key: Option, + #[serde(rename = "@title")] + pub title: Option, + #[serde(rename = "@type")] + pub section_type: Option, +} + +#[derive(Debug, Clone, Deserialize, PartialEq)] +#[serde(rename = "MediaContainer")] +pub struct PlexMediaContainerDto { + #[serde(rename = "@size")] + pub size: Option, + #[serde(rename = "@totalSize")] + pub total_size: Option, + #[serde(rename = "Video", default)] + pub videos: Vec, + #[serde(rename = "Directory", default)] + pub directories: Vec, +} + +#[derive(Debug, Clone, Deserialize, PartialEq)] +pub struct PlexDirectoryDto { + #[serde(rename = "@ratingKey")] + pub rating_key: Option, + #[serde(rename = "@key")] + pub key: Option, + #[serde(rename = "@type")] + pub item_type: Option, + #[serde(rename = "@title")] + pub title: Option, + #[serde(rename = "@year")] + pub year: Option, + #[serde(rename = "Guid", default)] + pub guids: Vec, +} + +#[derive(Debug, Clone, Deserialize, PartialEq)] +pub struct PlexVideoDto { + #[serde(rename = "@ratingKey")] + pub rating_key: Option, + #[serde(rename = "@key")] + pub key: Option, + #[serde(rename = "@type")] + pub item_type: Option, + #[serde(rename = "@title")] + pub title: Option, + #[serde(rename = "@year")] + pub year: Option, + #[serde(rename = "@guid")] + pub guid: Option, + #[serde(rename = "@thumb")] + pub thumb: Option, + #[serde(rename = "@art")] + pub art: Option, + #[serde(rename = "@parentRatingKey")] + pub parent_rating_key: Option, + #[serde(rename = "@grandparentRatingKey")] + pub grandparent_rating_key: Option, + #[serde(rename = "@grandparentTitle")] + pub grandparent_title: Option, + #[serde(rename = "@parentIndex")] + pub parent_index: Option, + #[serde(rename = "@index")] + pub index: Option, + #[serde(rename = "@addedAt")] + pub added_at: Option, + #[serde(rename = "@updatedAt")] + pub updated_at: Option, + #[serde(rename = "Guid", default)] + pub guids: Vec, + #[serde(rename = "Media", default)] + pub media: Vec, +} + +#[derive(Debug, Clone, Deserialize, PartialEq, Eq)] +pub struct PlexGuidDto { + #[serde(rename = "@id")] + pub id: Option, +} + +#[derive(Debug, Clone, Deserialize, PartialEq)] +pub struct PlexMediaDto { + #[serde(rename = "@id")] + pub id: Option, + #[serde(rename = "@container")] + pub container: Option, + #[serde(rename = "@duration")] + pub duration: Option, + #[serde(rename = "@bitrate")] + pub bitrate: Option, + #[serde(rename = "@width")] + pub width: Option, + #[serde(rename = "@height")] + pub height: Option, + #[serde(rename = "Part", default)] + pub parts: Vec, +} + +#[derive(Debug, Clone, Deserialize, PartialEq, Eq)] +pub struct PlexPartDto { + #[serde(rename = "@id")] + pub id: Option, + #[serde(rename = "@key")] + pub key: Option, + #[serde(rename = "@size")] + pub size: Option, + #[serde(rename = "@file")] + pub file: Option, + #[serde(rename = "@container")] + pub container: Option, +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::media_server::test_fixtures::{PLEX_MOVIES_XML, PLEX_RESOURCES_XML, PLEX_SECTIONS_XML}; + + #[test] + fn parses_plex_resources_without_exposing_resource_token() { + let resources: PlexResourcesDto = quick_xml::de::from_str(PLEX_RESOURCES_XML).expect("fixture parses"); + + assert_eq!(resources.devices.len(), 1); + assert_eq!(resources.devices[0].connections.len(), 1); + assert_eq!(resources.devices[0].access_token.as_deref(), Some("resource-token-redacted")); + } + + #[test] + fn parses_plex_sections_with_unsupported_kind_visible() { + let sections: PlexSectionsDto = quick_xml::de::from_str(PLEX_SECTIONS_XML).expect("fixture parses"); + + assert_eq!(sections.directories.len(), 3); + assert!(sections.directories.iter().any(|section| section.section_type.as_deref() == Some("artist"))); + } + + #[test] + fn parses_plex_movie_parts_and_optional_guids() { + let container: PlexMediaContainerDto = quick_xml::de::from_str(PLEX_MOVIES_XML).expect("fixture parses"); + let movie = &container.videos[0]; + let part = &movie.media[0].parts[0]; + + assert_eq!(container.total_size, Some(1)); + assert_eq!(movie.rating_key.as_deref(), Some("rating-redacted-1")); + assert_eq!(part.key.as_deref(), Some("/library/parts/part-redacted/file.mkv")); + assert!(part.file.as_deref().is_some_and(|file| file.contains("/redacted/"))); + assert_eq!(movie.guids.len(), 1); + } +} diff --git a/backend/src/media_server/plex/mod.rs b/backend/src/media_server/plex/mod.rs new file mode 100644 index 000000000..a07dce5c0 --- /dev/null +++ b/backend/src/media_server/plex/mod.rs @@ -0,0 +1 @@ +pub mod dto; diff --git a/backend/src/media_server/redaction.rs b/backend/src/media_server/redaction.rs new file mode 100644 index 000000000..0046b7b56 --- /dev/null +++ b/backend/src/media_server/redaction.rs @@ -0,0 +1,170 @@ +use reqwest::header::{HeaderMap, HeaderName}; +use shared::utils::sanitize_sensitive_info; + +const REDACTED_VALUE: &str = "***"; + +const SENSITIVE_QUERY_KEYS: &[&str] = &[ + "token", + "access_token", + "x-plex-token", + "x_emby_token", + "x-emby-token", + "x-mediabrowser-token", + "api_key", + "apikey", + "password", + "passwd", + "authorization", + "auth", +]; + +const SENSITIVE_HEADER_NAMES: &[&str] = &[ + "authorization", + "x-emby-token", + "x-mediabrowser-token", + "x-plex-token", + "cookie", + "set-cookie", +]; + +pub fn is_sensitive_media_server_header(name: &HeaderName) -> bool { + let value = name.as_str().to_ascii_lowercase(); + SENSITIVE_HEADER_NAMES.iter().any(|sensitive| value == *sensitive) +} + +pub fn redact_media_server_headers(headers: &HeaderMap) -> Vec<(String, String)> { + headers + .iter() + .map(|(name, value)| { + let rendered = if is_sensitive_media_server_header(name) { + REDACTED_VALUE.to_string() + } else { + value.to_str().map_or_else(|_| "".to_string(), redact_media_server_text) + }; + (name.as_str().to_string(), rendered) + }) + .collect() +} + +pub fn redact_media_server_text(value: &str) -> String { + let sanitized = sanitize_sensitive_info(value); + redact_query_like_tokens(&sanitized) +} + +fn redact_query_like_tokens(value: &str) -> String { + let mut result = String::with_capacity(value.len()); + let mut i = 0; + + while i < value.len() { + let matched = SENSITIVE_QUERY_KEYS.iter().find_map(|key| matched_sensitive_key(value, i, key)); + + if let Some((key, separator)) = matched { + result.push_str(key); + result.push(separator); + result.push_str(REDACTED_VALUE); + i += key.len() + separator.len_utf8(); + + let quote = value[i..].chars().next().filter(|ch| matches!(ch, '"' | '\'')); + if let Some(quote) = quote { + i += quote.len_utf8(); + while i < value.len() { + let Some(ch) = value[i..].chars().next() else { break }; + i += ch.len_utf8(); + if ch == quote { + break; + } + } + } else { + while i < value.len() { + let Some(ch) = value[i..].chars().next() else { break }; + if matches!(ch, '&' | ' ' | '\n' | '\r' | '\t' | '"' | '\'') { + break; + } + i += ch.len_utf8(); + } + } + } else { + let Some(ch) = value[i..].chars().next() else { break }; + result.push(ch); + i += ch.len_utf8(); + } + } + + result +} + +fn matched_sensitive_key(value: &str, start: usize, key: &'static str) -> Option<(&'static str, char)> { + let remaining = value.get(start..)?; + let candidate = remaining.get(..key.len())?; + let separator = remaining.get(key.len()..)?.chars().next()?; + if candidate.eq_ignore_ascii_case(key) && matches!(separator, '=' | ':') { + Some((key, separator)) + } else { + None + } +} + +#[cfg(test)] +mod tests { + use super::*; + use reqwest::header::{HeaderValue, AUTHORIZATION}; + + #[test] + fn redacts_media_server_query_tokens_case_insensitively() { + let redacted = redact_media_server_text( + "https://media.example.invalid/video?X-Plex-Token=secret-token&api_key=secret-key&safe=value", + ); + + assert!(!redacted.contains("secret-token")); + assert!(!redacted.contains("secret-key")); + assert!(redacted.contains("X-Plex-Token=***") || redacted.contains("x-plex-token=***")); + assert!(redacted.contains("api_key=***")); + assert!(redacted.contains("safe=value")); + } + + #[test] + fn redacts_query_tokens_without_corrupting_non_ascii_text() { + let redacted = redact_media_server_text("https://media.example.invalid/épisode?token=sëcret&title=café"); + + assert!(!redacted.contains("sëcret")); + assert!(redacted.contains("épisode")); + assert!(redacted.contains("title=café")); + assert!(redacted.contains("token=***")); + } + + #[test] + fn redacts_quoted_query_like_tokens() { + let redacted = redact_media_server_text( + "https://media.example.invalid/video?token=\"secret\"&api_key='another-secret'&safe=visible", + ); + + assert!(!redacted.contains("secret")); + assert!(!redacted.contains("another-secret")); + assert!(redacted.contains("token=***")); + assert!(redacted.contains("api_key=***")); + assert!(redacted.contains("safe=visible")); + } + + #[test] + fn redacts_media_browser_query_token() { + let redacted = redact_media_server_text( + "https://media.example.invalid/video?X-MediaBrowser-Token=secret-token&safe=visible", + ); + + assert!(!redacted.contains("secret-token")); + assert!(redacted.contains("X-MediaBrowser-Token=***") || redacted.contains("x-mediabrowser-token=***")); + assert!(redacted.contains("safe=visible")); + } + + #[test] + fn redacts_sensitive_headers() { + let mut headers = HeaderMap::new(); + headers.insert(AUTHORIZATION, HeaderValue::from_static("Bearer secret")); + headers.insert("content-type", HeaderValue::from_static("video/mp4")); + + let rendered = redact_media_server_headers(&headers); + + assert!(rendered.iter().any(|(name, value)| name == "authorization" && value == "***")); + assert!(rendered.iter().any(|(name, value)| name == "content-type" && value == "video/mp4")); + } +} diff --git a/backend/src/media_server/test_fixtures.rs b/backend/src/media_server/test_fixtures.rs new file mode 100644 index 000000000..111a516b9 --- /dev/null +++ b/backend/src/media_server/test_fixtures.rs @@ -0,0 +1,65 @@ +pub const EMBY_ITEMS_PAGE_JSON: &str = r#" +{ + "Items": [ + { + "Id": "item-redacted-1", + "Name": "Movie Redacted", + "Type": "Movie", + "ProductionYear": 2024, + "ProviderIds": { "Tmdb": "12345" }, + "ImageTags": { "Primary": "image-tag-redacted" }, + "Path": "/redacted/upstream/path/movie.mkv", + "UserData": { "Played": true, "PlaybackPositionTicks": 123 }, + "MediaSources": [ + { + "Id": "media-source-redacted-1", + "Container": "mkv", + "Path": "/redacted/upstream/path/movie.mkv", + "SupportsDirectPlay": true, + "SupportsDirectStream": true + } + ] + } + ], + "TotalRecordCount": 1, + "StartIndex": 0 +} +"#; + +pub const JELLYFIN_VIEWS_JSON: &str = r#" +{ + "Items": [ + { "Id": "library-redacted-movies", "Name": "Movies", "CollectionType": "movies", "Type": "CollectionFolder" }, + { "Id": "library-redacted-tv", "Name": "TV", "CollectionType": "tvshows", "Type": "CollectionFolder" } + ], + "TotalRecordCount": 2, + "StartIndex": 0 +} +"#; + +pub const PLEX_RESOURCES_XML: &str = r#" + + + + + +"#; + +pub const PLEX_SECTIONS_XML: &str = r#" + + + + + +"#; + +pub const PLEX_MOVIES_XML: &str = r#" + + + +"#; diff --git a/backend/src/media_server/types.rs b/backend/src/media_server/types.rs new file mode 100644 index 000000000..e4a74ce39 --- /dev/null +++ b/backend/src/media_server/types.rs @@ -0,0 +1,269 @@ +use crate::media_server::MediaServerError; +use bytes::Bytes; +use futures::stream::BoxStream; +use reqwest::{header::HeaderMap, StatusCode}; +use shared::model::InputType; +use std::{fmt, sync::Arc}; + +#[derive(Debug, Copy, Clone, PartialEq, Eq, Hash)] +pub enum MediaServerKind { + Emby, + Jellyfin, + Plex, +} + +impl MediaServerKind { + pub const fn as_input_type(self) -> InputType { + match self { + Self::Emby => InputType::Emby, + Self::Jellyfin => InputType::Jellyfin, + Self::Plex => InputType::Plex, + } + } +} + +impl TryFrom for MediaServerKind { + type Error = &'static str; + + fn try_from(value: InputType) -> Result { + match value { + InputType::Emby => Ok(Self::Emby), + InputType::Jellyfin => Ok(Self::Jellyfin), + InputType::Plex => Ok(Self::Plex), + InputType::M3u | InputType::Xtream | InputType::M3uBatch | InputType::XtreamBatch | InputType::Library => { + Err("input type is not a media-server input") + } + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MediaServerStatus { + pub kind: MediaServerKind, + pub server_id: Arc, + pub display_name: Option>, + pub version: Option>, + pub owned: Option, +} + +#[derive(Debug, Copy, Clone, PartialEq, Eq, Hash)] +pub enum MediaServerLibraryKind { + Movies, + TvShows, + Unsupported, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MediaServerLibraryRef { + pub input_name: Arc, + pub server_id: Arc, + pub library_id: Arc, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MediaServerLibrary { + pub reference: MediaServerLibraryRef, + pub name: Arc, + pub kind: MediaServerLibraryKind, +} + +#[derive(Debug, Copy, Clone, PartialEq, Eq)] +pub struct MediaServerPageRequest { + pub start: usize, + pub limit: usize, +} + +impl MediaServerPageRequest { + pub const fn new(start: usize, limit: usize) -> Self { Self { start, limit } } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MediaServerPage { + pub request: MediaServerPageRequest, + pub total: Option, + pub upstream_item_count: usize, + pub items: Vec, +} + +impl MediaServerPage { + pub fn new(request: MediaServerPageRequest, total: Option, items: Vec) -> Self { + let upstream_item_count = items.len(); + Self { request, total, upstream_item_count, items } + } + + pub fn with_upstream_item_count( + request: MediaServerPageRequest, + total: Option, + upstream_item_count: usize, + items: Vec, + ) -> Self { + debug_assert!( + upstream_item_count >= items.len(), + "upstream_item_count must be greater than or equal to items.len()" + ); + let upstream_item_count = upstream_item_count.max(items.len()); + Self { request, total, upstream_item_count, items } + } + + pub fn item_count(&self) -> usize { self.items.len() } + + pub fn upstream_item_count(&self) -> usize { self.upstream_item_count } + + pub fn next_request(&self) -> Option { + let next_start = self.request.start.saturating_add(self.upstream_item_count()); + if self.upstream_item_count() == 0 || self.total.is_some_and(|total| next_start >= total) { + None + } else { + Some(MediaServerPageRequest::new(next_start, self.request.limit)) + } + } + + pub fn cursor_advanced(&self) -> bool { self.upstream_item_count() > 0 } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MediaServerProviderIdHint { + pub namespace: Arc, + pub value: Arc, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MediaServerMovie { + pub input_name: Arc, + pub server_id: Arc, + pub library_id: Arc, + pub item_id: Arc, + pub title: Arc, + pub year: Option, + pub source_version_hint: Option>, + pub provider_hints: Vec, + pub stream_ref: Option, + pub image_ref: Option, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MediaServerEpisode { + pub input_name: Arc, + pub server_id: Arc, + pub library_id: Arc, + pub item_id: Arc, + pub series_id: Option>, + pub series_title: Option>, + pub title: Arc, + pub season: Option, + pub episode: Option, + pub source_version_hint: Option>, + pub provider_hints: Vec, + pub stream_ref: Option, + pub image_ref: Option, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum MediaServerStreamRef { + Emby { + input_name: Arc, + server_id: Arc, + item_id: Arc, + media_source_id: Option>, + }, + Jellyfin { + input_name: Arc, + server_id: Arc, + item_id: Arc, + media_source_id: Option>, + }, + Plex { + input_name: Arc, + server_id: Arc, + rating_key: Arc, + part_key: Arc, + }, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum MediaServerImageRef { + Emby { + input_name: Arc, + server_id: Arc, + item_id: Arc, + image_kind: Arc, + tag: Option>, + }, + Jellyfin { + input_name: Arc, + server_id: Arc, + item_id: Arc, + image_kind: Arc, + tag: Option>, + }, + Plex { + input_name: Arc, + server_id: Arc, + rating_key: Arc, + image_path: Arc, + }, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MediaServerPlaybackLease { + pub provider_kind: MediaServerKind, + pub lease_id: Arc, +} + +#[derive(Debug, Clone)] +pub struct MediaServerResourceResponse { + pub status: StatusCode, + pub headers: HeaderMap, + pub body: Bytes, +} + +pub type BoxedMediaServerStream = BoxStream<'static, Result>; + +pub struct MediaServerStreamResponse { + pub status: StatusCode, + pub headers: HeaderMap, + pub body: BoxedMediaServerStream, +} + +impl fmt::Debug for MediaServerStreamResponse { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("MediaServerStreamResponse") + .field("status", &self.status) + .field("headers", &self.headers) + .field("body", &"") + .finish() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn media_server_page_next_request_advances_until_total() { + let page = MediaServerPage::new(MediaServerPageRequest::new(0, 2), Some(3), vec![1, 2]); + + assert_eq!(page.next_request(), Some(MediaServerPageRequest::new(2, 2))); + + let last = MediaServerPage::new(MediaServerPageRequest::new(2, 2), Some(3), vec![3]); + assert_eq!(last.next_request(), None); + } + + #[test] + fn media_server_page_empty_page_does_not_advance() { + let page = MediaServerPage::::new(MediaServerPageRequest::new(10, 100), Some(50), vec![]); + + assert!(!page.cursor_advanced()); + assert_eq!(page.next_request(), None); + } + + #[test] + fn media_server_page_cursor_uses_upstream_count_when_items_are_filtered() { + let page = MediaServerPage::with_upstream_item_count(MediaServerPageRequest::new(0, 3), Some(5), 3, vec![1]); + + assert_eq!(page.item_count(), 1); + assert_eq!(page.upstream_item_count(), 3); + assert!(page.cursor_advanced()); + assert_eq!(page.next_request(), Some(MediaServerPageRequest::new(3, 3))); + } +} diff --git a/backend/src/model/config/input.rs b/backend/src/model/config/input.rs index 1c1b72098..4cb61695d 100644 --- a/backend/src/model/config/input.rs +++ b/backend/src/model/config/input.rs @@ -5,8 +5,15 @@ use log::warn; use shared::foundation::Filter; use shared::{apply_flags, create_bitset}; use shared::error::TuliproxError; -use shared::model::{ClusterSource, ConfigInputAliasDto, ConfigInputDto, ConfigInputOptionsDto, InputFetchMethod, InputType, StagedInputDto, XtreamCluster}; -use shared::utils::{get_credentials_from_url, parse_provider_scheme_url_parts, sanitize_sensitive_info, Internable, PROVIDER_SCHEME_PREFIX}; +use shared::model::{ + ClusterSource, ConfigInputAliasDto, ConfigInputDto, ConfigInputOptionsDto, InputFetchMethod, InputType, + MediaServerCatalogConfigDto, MediaServerEnrichmentConfigDto, MediaServerImagePolicyDto, MediaServerLibrarySelectorDto, + MediaServerInputConfigDto, MediaServerPlaybackConfigDto, StagedInputDto, XtreamCluster, +}; +use shared::utils::{ + get_credentials_from_url, is_non_blank_optional_string, parse_provider_scheme_url_parts, sanitize_sensitive_info, Internable, + BATCH_SCHEME_PREFIX, PROVIDER_SCHEME_PREFIX, +}; use shared::{check_input_connections, write_if_some}; use shared::{check_input_credentials, concat_string }; use std::borrow::Cow; @@ -99,6 +106,63 @@ impl From<&ConfigInputOptionsDto> for ConfigInputOptions { static DEFAULT_CONFIG_INPUT_OPTIONS: LazyLock = LazyLock::new(|| ConfigInputOptions::from(&ConfigInputOptionsDto::default())); +#[derive(Debug, Clone)] +pub struct MediaServerInputConfig { + pub libraries: Vec, + pub catalog: MediaServerCatalogConfigDto, + pub playback: MediaServerPlaybackConfigDto, + pub enrichment: MediaServerEnrichmentConfigDto, + pub image_policy: MediaServerImagePolicyDto, + pub token: Option, + pub api_key: Option, + pub user_id: Option, + pub account_token: Option, + pub server_id: Option, + pub machine_id: Option, + pub server_name: Option, + pub prefer_https: bool, + pub allow_relay: bool, +} + +impl From<&MediaServerInputConfigDto> for MediaServerInputConfig { + fn from(dto: &MediaServerInputConfigDto) -> Self { + let mut normalized = dto.clone(); + normalized.normalize(); + Self { + libraries: normalized.libraries, + catalog: normalized.catalog, + playback: normalized.playback, + enrichment: normalized.enrichment, + image_policy: normalized.image_policy, + token: normalized.token, + api_key: normalized.api_key, + user_id: normalized.user_id, + account_token: normalized.account_token, + server_id: normalized.server_id, + machine_id: normalized.machine_id, + server_name: normalized.server_name, + prefer_https: normalized.prefer_https, + allow_relay: normalized.allow_relay, + } + } +} + +impl MediaServerInputConfig { + pub fn has_any_emby_jellyfin_auth(&self) -> bool { + is_non_blank_optional_string(&self.token) || is_non_blank_optional_string(&self.api_key) + } + + pub fn has_any_plex_token(&self) -> bool { + is_non_blank_optional_string(&self.account_token) || is_non_blank_optional_string(&self.token) + } + + pub fn has_plex_server_selector(&self) -> bool { + is_non_blank_optional_string(&self.server_id) + || is_non_blank_optional_string(&self.machine_id) + || is_non_blank_optional_string(&self.server_name) + } +} + pub struct InputUserInfo { pub base_url: String, pub username: String, @@ -254,6 +318,7 @@ pub struct ConfigInput { pub persist: Option, pub enabled: bool, pub options: Option, + pub media_server: Option, pub aliases: Option>, pub priority: i16, pub max_connections: u16, @@ -434,6 +499,103 @@ impl ConfigInput { self.options.as_ref().map_or(default, |o| o.has_all_flags(flags)) } + fn prepare_media_server_input(&self) -> Result<(), TuliproxError> { + if !self.input_type.is_media_server() { + return Ok(()); + } + + let trimmed_url = self.url.trim(); + if trimmed_url.starts_with(BATCH_SCHEME_PREFIX) || trimmed_url.starts_with(PROVIDER_SCHEME_PREFIX) { + return Err(TuliproxError::ConfigInput(format!( + "media-server input does not support batch:// or provider:// URLs (input: {})", + self.name + ))); + } + if self.aliases.as_ref().is_some_and(|aliases| !aliases.is_empty()) { + return Err(TuliproxError::ConfigInput(format!( + "media-server input does not support aliases (input: {})", + self.name + ))); + } + if self.staged.as_ref().is_some_and(|staged| staged.enabled) { + return Err(TuliproxError::ConfigInput(format!( + "media-server input does not support staged inputs (input: {})", + self.name + ))); + } + if self.epg.is_some() { + return Err(TuliproxError::ConfigInput(format!( + "media-server input does not support EPG configuration (input: {})", + self.name + ))); + } + if self.panel_api.is_some() { + return Err(TuliproxError::ConfigInput(format!( + "media-server input does not support panel_api configuration (input: {})", + self.name + ))); + } + let Some(media_server) = self.media_server.as_ref() else { + return Err(TuliproxError::ConfigInput(format!( + "media_server configuration is mandatory for input type {} (input: {})", + self.input_type, self.name + ))); + }; + if media_server.libraries.is_empty() { + return Err(TuliproxError::ConfigInput(format!( + "media-server input requires at least one selected library (input: {})", + self.name + ))); + } + if media_server.libraries.iter().any(MediaServerLibrarySelectorDto::is_empty) { + return Err(TuliproxError::ConfigInput(format!( + "media_server library selectors must not be empty (input: {})", + self.name + ))); + } + if media_server.catalog.page_size == 0 { + return Err(TuliproxError::ConfigInput(format!( + "media server catalog page_size must be greater than zero (input: {})", + self.name + ))); + } + match self.input_type { + InputType::Emby | InputType::Jellyfin => { + if trimmed_url.is_empty() { + return Err(TuliproxError::ConfigInput(format!( + "url is mandatory for input type {} (input: {})", + self.input_type, self.name + ))); + } + let has_login = self.username.as_ref().is_some_and(|u| !u.trim().is_empty()) + && self.password.as_ref().is_some_and(|p| !p.trim().is_empty()); + if !media_server.has_any_emby_jellyfin_auth() && !has_login { + return Err(TuliproxError::ConfigInput(format!( + "media-server input type {} requires media_server token/api_key or username/password bootstrap credentials (input: {})", + self.input_type, self.name + ))); + } + } + InputType::Plex => { + if !media_server.has_any_plex_token() { + return Err(TuliproxError::ConfigInput(format!( + "media-server input type plex requires media_server.account_token or media_server.token (input: {})", + self.name + ))); + } + if trimmed_url.is_empty() && !media_server.has_plex_server_selector() { + return Err(TuliproxError::ConfigInput(format!( + "media-server input type plex requires a server selector such as media_server.machine_id, media_server.server_id, or media_server.server_name when input.url is omitted (input: {})", + self.name + ))); + } + } + InputType::M3u | InputType::Xtream | InputType::M3uBatch | InputType::XtreamBatch | InputType::Library => {} + } + + Ok(()) + } + pub fn prepare(&mut self, provider_configs: &[Arc]) -> Result, TuliproxError> { // Defensive fallback: From<&ConfigInputDto> for ConfigInput sets options, but ConfigInput can // still be built via Default::default(), batch/internal/test paths, so prepare() normalizes @@ -452,6 +614,10 @@ impl ConfigInput { let batch_file_path = self.prepare_batch(); self.name = self.name.trim().intern(); + if self.enabled { + self.prepare_media_server_input()?; + } + if self.url.starts_with(PROVIDER_SCHEME_PREFIX) { let provider_cfg = Self::resolve_provider_config(&self.url, provider_configs)?; used_provider_configs.push(provider_cfg); @@ -557,6 +723,7 @@ impl ConfigInput { persist: self.persist.clone(), enabled: self.enabled, options: self.options.clone(), + media_server: self.media_server.clone(), aliases: None, priority: alias.priority, max_connections: alias.max_connections, @@ -644,6 +811,7 @@ impl From<&ConfigInputDto> for ConfigInput { persist: dto.persist.clone(), enabled: dto.enabled, options: Some(options), + media_server: dto.media_server.as_ref().map(MediaServerInputConfig::from), aliases: dto.aliases.as_ref().map(|list| list.iter().map(ConfigInputAlias::from).collect()), priority: dto.priority, max_connections: dto.max_connections, @@ -755,6 +923,190 @@ mod tests { use std::borrow::Cow; use std::sync::Arc; + fn media_server_config_with_library() -> MediaServerInputConfig { + MediaServerInputConfig::from(&MediaServerInputConfigDto { + libraries: vec![MediaServerLibrarySelectorDto::Name("Movies".to_string())], + ..MediaServerInputConfigDto::default() + }) + } + + #[test] + fn media_server_runtime_mapping_preserves_safe_defaults() { + let dto = ConfigInputDto { + name: "emby_media_server".into(), + input_type: InputType::Emby, + url: "https://media.example.invalid".to_string(), + media_server: Some(MediaServerInputConfigDto { + token: Some(" token ".to_string()), + api_key: Some(" api-key ".to_string()), + user_id: Some(" user ".to_string()), + account_token: Some(" account-token ".to_string()), + server_id: Some(" server ".to_string()), + machine_id: Some(" machine ".to_string()), + server_name: Some(" server-name ".to_string()), + ..MediaServerInputConfigDto { + libraries: vec![MediaServerLibrarySelectorDto::Name(" Movies ".to_string())], + ..MediaServerInputConfigDto::default() + } + }), + ..ConfigInputDto::default() + }; + + let input = ConfigInput::from(&dto); + let media_server = input.media_server.expect("media_server config should map to runtime"); + + assert_eq!(media_server.libraries, vec![MediaServerLibrarySelectorDto::Name("Movies".to_string())]); + assert_eq!(media_server.catalog.page_size, 100); + assert!(media_server.playback.direct_play_only); + assert!(!media_server.playback.allow_transcode); + assert!(!media_server.enrichment.ffprobe); + assert_eq!(media_server.token.as_deref(), Some("token")); + assert_eq!(media_server.api_key.as_deref(), Some("api-key")); + assert_eq!(media_server.user_id.as_deref(), Some("user")); + assert_eq!(media_server.account_token.as_deref(), Some("account-token")); + assert_eq!(media_server.server_id.as_deref(), Some("server")); + assert_eq!(media_server.machine_id.as_deref(), Some("machine")); + assert_eq!(media_server.server_name.as_deref(), Some("server-name")); + } + + #[test] + fn prepare_accepts_plex_media_server_without_input_url() { + let mut input = ConfigInput { + name: "plex_media_server".into(), + input_type: InputType::Plex, + media_server: Some(MediaServerInputConfig { + account_token: Some("token".to_string()), + machine_id: Some("machine".to_string()), + ..media_server_config_with_library() + }), + enabled: true, + ..Default::default() + }; + + input.prepare(&[]).expect("plex discovery config should prepare without input.url"); + } + + #[test] + fn prepare_accepts_plex_media_server_with_direct_url_without_selector() { + let mut input = ConfigInput { + name: "plex_media_server".into(), + input_type: InputType::Plex, + url: "https://plex.example.invalid".to_string(), + media_server: Some(MediaServerInputConfig { + token: Some("token".to_string()), + ..media_server_config_with_library() + }), + enabled: true, + ..Default::default() + }; + + input.prepare(&[]).expect("direct Plex URL should not require MyPlex server selector"); + } + + #[test] + fn prepare_rejects_emby_media_server_without_input_url() { + let mut input = ConfigInput { + name: "emby_media_server".into(), + input_type: InputType::Emby, + media_server: Some(MediaServerInputConfig { + token: Some("token".to_string()), + ..media_server_config_with_library() + }), + enabled: true, + ..Default::default() + }; + + let err = input.prepare(&[]).expect_err("emby media_server input should require a direct server URL"); + assert!(err.to_string().contains("url is mandatory for input type emby")); + } + + #[test] + fn prepare_rejects_blank_media_server_credentials_and_selectors() { + let mut emby = ConfigInput { + name: "emby_media_server".into(), + input_type: InputType::Emby, + url: "https://media.example.invalid".to_string(), + media_server: Some(MediaServerInputConfig { + token: Some(" ".to_string()), + api_key: Some(String::new()), + ..media_server_config_with_library() + }), + enabled: true, + ..Default::default() + }; + let err = emby.prepare(&[]).expect_err("blank token/api_key should be rejected"); + assert!(err.to_string().contains("requires media_server token/api_key")); + + let mut plex = ConfigInput { + name: "plex_media_server".into(), + input_type: InputType::Plex, + media_server: Some(MediaServerInputConfig { + account_token: Some(" ".to_string()), + server_id: Some(" ".to_string()), + ..media_server_config_with_library() + }), + enabled: true, + ..Default::default() + }; + let err = plex.prepare(&[]).expect_err("blank plex token should be rejected"); + assert!(err.to_string().contains("requires media_server.account_token or media_server.token")); + } + + #[test] + fn prepare_rejects_blank_media_server_library_selector() { + let mut input = ConfigInput { + name: "emby_media_server".into(), + input_type: InputType::Emby, + url: "https://media.example.invalid".to_string(), + media_server: Some(MediaServerInputConfig { + token: Some("token".to_string()), + libraries: vec![MediaServerLibrarySelectorDto::Name(" ".to_string())], + ..media_server_config_with_library() + }), + enabled: true, + ..Default::default() + }; + + let err = input.prepare(&[]).expect_err("blank library selector should be rejected"); + assert!(err.to_string().contains("media_server library selectors must not be empty")); + } + + #[test] + fn prepare_accepts_media_server_max_connections_as_stream_limit() { + let mut input = ConfigInput { + name: "emby_media_server".into(), + input_type: InputType::Emby, + url: "https://media.example.invalid".to_string(), + media_server: Some(MediaServerInputConfig { + token: Some("token".to_string()), + ..media_server_config_with_library() + }), + max_connections: 1, + enabled: true, + ..Default::default() + }; + + input.prepare(&[]).expect("media_server inputs reuse max_connections stream-limit semantics"); + } + + #[test] + fn prepare_rejects_media_server_provider_url() { + let mut input = ConfigInput { + name: "emby_media_server".into(), + input_type: InputType::Emby, + url: " provider://media-server ".to_string(), + media_server: Some(MediaServerInputConfig { + token: Some("token".to_string()), + ..media_server_config_with_library() + }), + enabled: true, + ..Default::default() + }; + + let err = input.prepare(&[]).expect_err("media_server provider URLs should be rejected"); + assert!(err.to_string().contains("does not support batch:// or provider://")); + } + #[test] fn test_resolve_url_normal() { let input = ConfigInput { diff --git a/backend/src/modules.rs b/backend/src/modules.rs index 6aba7a01e..7b49c7f2a 100644 --- a/backend/src/modules.rs +++ b/backend/src/modules.rs @@ -7,6 +7,7 @@ macro_rules! include_modules { pub mod api; pub mod auth; pub mod library; + pub mod media_server; pub mod messaging; pub mod model; pub mod processing; diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 06eb97a35..3c9f3219a 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -507,6 +507,16 @@ async fn playlist_download_from_input( let (p, e) = library::download_library_playlist(client, app_config, input).await; (p, e, false, 0, 0) } + InputType::Emby | InputType::Jellyfin | InputType::Plex => ( + vec![], + vec![TuliproxError::Download(format!( + "media-server input '{}' is configured but catalog import is not implemented yet", + input.name + ))], + false, + 0, + 0, + ), } }; diff --git a/backend/src/processing/processor/stream_probe.rs b/backend/src/processing/processor/stream_probe.rs index 6261da1a6..44e0b16c2 100644 --- a/backend/src/processing/processor/stream_probe.rs +++ b/backend/src/processing/processor/stream_probe.rs @@ -58,7 +58,7 @@ pub enum GenericProbeMetadataOutcome { } fn requires_provider_connection_for_generic_probe(input_type: InputType) -> bool { - !matches!(input_type, InputType::Library) + !(matches!(input_type, InputType::Library) || input_type.is_media_server()) } fn uses_seekable_remote_probe(item_type: PlaylistItemType, is_remote_probe: bool) -> bool { @@ -183,13 +183,14 @@ async fn prepare_generic_stream_metadata( XtreamCluster::Series } else { // Generic probing currently supports live/video/series payload shapes. - return Ok(PreparedGenericProbeOutcome::Noop); + return Ok(PreparedGenericProbeOutcome::Noop); }; ( xtream_get_file_path(&storage_path, cluster), ProbeStorageKind::Xtream, ) } + InputType::Emby | InputType::Jellyfin | InputType::Plex => return Ok(PreparedGenericProbeOutcome::Noop), }; if !db_path.exists() { @@ -532,6 +533,13 @@ mod tests { assert!(!requires_provider_connection_for_generic_probe(InputType::Library)); } + #[test] + fn media_server_probe_does_not_require_provider_connection() { + assert!(!requires_provider_connection_for_generic_probe(InputType::Emby)); + assert!(!requires_provider_connection_for_generic_probe(InputType::Jellyfin)); + assert!(!requires_provider_connection_for_generic_probe(InputType::Plex)); + } + #[test] fn m3u_probe_requires_provider_connection() { assert!(requires_provider_connection_for_generic_probe(InputType::M3u)); diff --git a/backend/src/repository/alias_repository.rs b/backend/src/repository/alias_repository.rs index 6cb1c80d2..1cc62cbf0 100644 --- a/backend/src/repository/alias_repository.rs +++ b/backend/src/repository/alias_repository.rs @@ -179,6 +179,7 @@ pub fn csv_read_inputs_from_reader( InputType::M3uBatch | InputType::M3u => InputType::M3uBatch, InputType::XtreamBatch | InputType::Xtream => InputType::XtreamBatch, InputType::Library => InputType::Library, + InputType::Emby | InputType::Jellyfin | InputType::Plex => batch_input_type, }; let mut result = vec![]; let mut default_columns = vec![]; diff --git a/backend/src/repository/playlist_repository.rs b/backend/src/repository/playlist_repository.rs index 42bcc31ba..adab1a65c 100644 --- a/backend/src/repository/playlist_repository.rs +++ b/backend/src/repository/playlist_repository.rs @@ -10,7 +10,10 @@ use crate::repository::{load_input_local_library_playlist, persist_input_library use crate::repository::{load_input_m3u_playlist, m3u_get_file_path_for_db, m3u_write_playlist, persist_input_m3u_playlist}; use crate::repository::{load_input_xtream_playlist, persist_input_xtream_playlist, xtream_get_file_path, xtream_get_storage_path, xtream_write_playlist}; use crate::repository::BPlusTree; -use crate::repository::{LocalLibraryDiskPlaylistSource, M3uDiskPlaylistSource, MemoryPlaylistSource, PlaylistSource, XtreamDiskPlaylistSource}; +use crate::repository::{ + LocalLibraryDiskPlaylistSource, M3uDiskPlaylistSource, MemoryPlaylistSource, PlaylistSource, + MediaServerDiskPlaylistSource, XtreamDiskPlaylistSource, +}; use crate::repository::{TargetIdMapping, VirtualIdRecord}; use crate::utils; use log::{info, warn}; @@ -459,6 +462,14 @@ pub async fn persist_input_playlist(app_config: &Arc, input: &ConfigI } (playlist, None) } + InputType::Emby | InputType::Jellyfin | InputType::Plex => { + let file_path = get_input_media_server_playlist_file_path(&storage_path, &input.name); + let (playlist, result) = persist_input_media_server_playlist(app_config, &file_path, playlist).await; + if let Err(err) = result { + return (playlist, Some(err)); + } + (playlist, None) + } } } @@ -502,6 +513,15 @@ pub async fn load_input_playlist(ctx: &PlaylistProcessingContext, input: &Config Ok(Box::new(MemoryPlaylistSource::new(groups))) } } + InputType::Emby | InputType::Jellyfin | InputType::Plex => { + let file_path = get_input_media_server_playlist_file_path(&storage_path, &input.name); + if disk_based_processing && file_path.exists() { + Ok(Box::new(MediaServerDiskPlaylistSource::new(app_config, &file_path).await)) + } else { + let groups = load_input_media_server_playlist(app_config, &file_path).await?; + Ok(Box::new(MemoryPlaylistSource::new(groups))) + } + } } } @@ -519,12 +539,34 @@ pub fn get_input_local_library_playlist_file_path(storage_path: &Path, input_nam storage_path.join(format!("lib_{sanitized_input_name}.{FILE_SUFFIX_DB}")) } +pub fn get_input_media_server_playlist_file_path(storage_path: &Path, input_name: &Arc) -> PathBuf { + let sanitized_input_name: String = input_name.chars() + .map(|c| if c.is_alphanumeric() { c } else { '_' }) + .collect(); + storage_path.join(format!("media_server_{sanitized_input_name}.{FILE_SUFFIX_DB}")) +} + +pub async fn persist_input_media_server_playlist( + app_config: &Arc, + file_path: &Path, + playlist: Vec, +) -> (Vec, Result<(), TuliproxError>) { + persist_input_library_playlist(app_config, file_path, playlist).await +} + +pub async fn load_input_media_server_playlist( + app_config: &Arc, + file_path: &Path, +) -> Result, TuliproxError> { + load_input_local_library_playlist(app_config, file_path).await +} + #[cfg(test)] mod tests { use super::{ - assign_local_series_info_episode_key, rewrite_local_series_info_episode_virtual_id, - rewrite_series_episode_parent_virtual_ids, rewrite_series_info_episode_virtual_id, LocalEpisodeKey, - ProviderEpisodeKey, + assign_local_series_info_episode_key, get_input_media_server_playlist_file_path, + rewrite_local_series_info_episode_virtual_id, rewrite_series_episode_parent_virtual_ids, + rewrite_series_info_episode_virtual_id, LocalEpisodeKey, ProviderEpisodeKey, }; use crate::repository::{BPlusTreeQuery, TargetIdMapping, VirtualIdRecord}; use shared::model::{ @@ -535,6 +577,14 @@ mod tests { use std::{collections::HashMap, sync::Arc}; use tempfile::tempdir; + #[test] + fn media_server_playlist_file_path_uses_separate_prefix() { + let dir = tempdir().expect("tempdir"); + let path = get_input_media_server_playlist_file_path(dir.path(), &"Media Server Input".intern()); + + assert!(path.ends_with("media_server_Media_Server_Input.db")); + } + fn make_local_series_info(series_uuid: &str, episodes: Vec<(u32, &str, &str)>) -> PlaylistItem { let episode_props = episodes .into_iter() diff --git a/backend/src/repository/playlist_source.rs b/backend/src/repository/playlist_source.rs index bf8273345..b27364eba 100644 --- a/backend/src/repository/playlist_source.rs +++ b/backend/src/repository/playlist_source.rs @@ -436,6 +436,7 @@ macro_rules! impl_single_file_disk_source { impl_single_file_disk_source!(M3u, Arc, M3uPlaylistItem); impl_single_file_disk_source!(LocalLibrary, UUIDType, XtreamPlaylistItem); +impl_single_file_disk_source!(MediaServer, UUIDType, XtreamPlaylistItem); pub struct MemoryPlaylistSource { playlist: Arc>, diff --git a/frontend/public/assets/i18n/en.json b/frontend/public/assets/i18n/en.json index 9fe62f885..379777ae5 100644 --- a/frontend/public/assets/i18n/en.json +++ b/frontend/public/assets/i18n/en.json @@ -1631,8 +1631,11 @@ } }, "SOURCE_EDITOR": { + "BRICK_InputEmby": "Emby", + "BRICK_InputJellyfin": "Jellyfin", "BRICK_InputLibrary": "Library", "BRICK_InputM3u": "M3u", + "BRICK_InputPlex": "Plex", "BRICK_InputXtream": "Xtream", "BRICK_OutputHdHomeRun": "HDHR", "BRICK_OutputM3u": "M3u", diff --git a/frontend/src/app/components/mod.rs b/frontend/src/app/components/mod.rs index 9ea104543..6aac59d13 100644 --- a/frontend/src/app/components/mod.rs +++ b/frontend/src/app/components/mod.rs @@ -53,15 +53,14 @@ mod userlist; mod websocket_status; mod cluster_flags_input; +mod country; mod field_explanation; mod field_id; mod filter; mod particle_flow_background; +mod setup; mod source_editor; mod title_card; - -mod country; -mod setup; // pub use self::input::*; // pub use self::menu_item::*; // pub use self::popup_menu::*; diff --git a/frontend/src/app/components/playlist/input/input_type_view.rs b/frontend/src/app/components/playlist/input/input_type_view.rs index 4f8f2c2e3..3c03088d5 100644 --- a/frontend/src/app/components/playlist/input/input_type_view.rs +++ b/frontend/src/app/components/playlist/input/input_type_view.rs @@ -17,6 +17,9 @@ pub fn InputTypeView(props: &InputTypeViewProps) -> Html { InputType::M3uBatch => "LABEL.M3U_BATCH", InputType::XtreamBatch => "LABEL.XTREAM_BATCH", InputType::Library => "LABEL.LIBRARY", + InputType::Emby => "LABEL.EMBY", + InputType::Jellyfin => "LABEL.JELLYFIN", + InputType::Plex => "LABEL.PLEX", }; html! { diff --git a/frontend/src/app/components/playlist/input/staged_input_view.rs b/frontend/src/app/components/playlist/input/staged_input_view.rs index 4e4056be6..bcb286fed 100644 --- a/frontend/src/app/components/playlist/input/staged_input_view.rs +++ b/frontend/src/app/components/playlist/input/staged_input_view.rs @@ -23,6 +23,9 @@ pub fn StagedInputView(props: &StagedInputViewProps) -> Html { InputType::M3uBatch => "LABEL.M3U_BATCH", InputType::XtreamBatch => "LABEL.XTREAM_BATCH", InputType::Library => "LABEL.LIBRARY", + InputType::Emby => "LABEL.EMBY", + InputType::Jellyfin => "LABEL.JELLYFIN", + InputType::Plex => "LABEL.PLEX", }; html! {
diff --git a/frontend/src/app/components/source_editor/editor_model.rs b/frontend/src/app/components/source_editor/editor_model.rs index 6e62f1a14..b6c9ab481 100644 --- a/frontend/src/app/components/source_editor/editor_model.rs +++ b/frontend/src/app/components/source_editor/editor_model.rs @@ -14,6 +14,9 @@ pub enum BlockType { InputXtream, InputM3u, InputLibrary, + InputEmby, + InputJellyfin, + InputPlex, Target, OutputM3u, OutputXtream, @@ -26,13 +29,26 @@ impl BlockType { pub const INPUT_XTREAM: &'static str = "InputXtream"; pub const INPUT_M3U: &'static str = "InputM3u"; pub const INPUT_LIBRARY: &'static str = "InputLibrary"; + pub const INPUT_EMBY: &'static str = "InputEmby"; + pub const INPUT_JELLYFIN: &'static str = "InputJellyfin"; + pub const INPUT_PLEX: &'static str = "InputPlex"; pub const TARGET: &'static str = "Target"; pub const OUTPUT_M3U: &'static str = "OutputM3u"; pub const OUTPUT_XTREAM: &'static str = "OutputXtream"; pub const OUTPUT_HDHOMERUN: &'static str = "OutputHdHomeRun"; pub const OUTPUT_STRM: &'static str = "OutputStrm"; - pub fn is_input(&self) -> bool { matches!(self, Self::InputXtream | Self::InputM3u | Self::InputLibrary) } + pub fn is_input(&self) -> bool { + matches!( + self, + Self::InputXtream + | Self::InputM3u + | Self::InputLibrary + | Self::InputEmby + | Self::InputJellyfin + | Self::InputPlex + ) + } pub fn is_target(&self) -> bool { matches!(self, Self::Target) } @@ -48,6 +64,9 @@ impl From<&str> for BlockType { BlockType::INPUT_XTREAM => BlockType::InputXtream, BlockType::INPUT_M3U => BlockType::InputM3u, BlockType::INPUT_LIBRARY => BlockType::InputLibrary, + BlockType::INPUT_EMBY => BlockType::InputEmby, + BlockType::INPUT_JELLYFIN => BlockType::InputJellyfin, + BlockType::INPUT_PLEX => BlockType::InputPlex, BlockType::TARGET => BlockType::Target, BlockType::OUTPUT_M3U => BlockType::OutputM3u, BlockType::OUTPUT_XTREAM => BlockType::OutputXtream, @@ -68,6 +87,9 @@ impl From for BlockType { InputType::M3uBatch | InputType::M3u => BlockType::InputM3u, InputType::XtreamBatch | InputType::Xtream => BlockType::InputXtream, InputType::Library => BlockType::InputLibrary, + InputType::Emby => BlockType::InputEmby, + InputType::Jellyfin => BlockType::InputJellyfin, + InputType::Plex => BlockType::InputPlex, } } } @@ -79,6 +101,9 @@ impl fmt::Display for BlockType { BlockType::InputXtream => Self::INPUT_XTREAM, BlockType::InputM3u => Self::INPUT_M3U, BlockType::InputLibrary => Self::INPUT_LIBRARY, + BlockType::InputEmby => Self::INPUT_EMBY, + BlockType::InputJellyfin => Self::INPUT_JELLYFIN, + BlockType::InputPlex => Self::INPUT_PLEX, BlockType::Target => Self::TARGET, BlockType::OutputM3u => Self::OUTPUT_M3U, BlockType::OutputXtream => Self::OUTPUT_XTREAM, @@ -90,6 +115,25 @@ impl fmt::Display for BlockType { } pub(crate) type BlockId = u16; +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn media_server_input_types_keep_distinct_source_editor_blocks() { + assert_eq!(BlockType::from(InputType::Emby), BlockType::InputEmby); + assert_eq!(BlockType::from(InputType::Jellyfin), BlockType::InputJellyfin); + assert_eq!(BlockType::from(InputType::Plex), BlockType::InputPlex); + } + + #[test] + fn media_server_source_editor_blocks_are_inputs() { + assert!(BlockType::InputEmby.is_input()); + assert!(BlockType::InputJellyfin.is_input()); + assert!(BlockType::InputPlex.is_input()); + } +} + #[derive(Clone, Debug, Serialize, Deserialize, PartialEq)] pub struct Block { pub id: BlockId, diff --git a/frontend/src/app/components/source_editor/editor_view.rs b/frontend/src/app/components/source_editor/editor_view.rs index 13d859fad..c69e46871 100644 --- a/frontend/src/app/components/source_editor/editor_view.rs +++ b/frontend/src/app/components/source_editor/editor_view.rs @@ -268,6 +268,9 @@ fn create_instance(block_type: BlockType) -> BlockInstance { BlockType::InputXtream => BlockInstance::Input(Rc::new(ConfigInputDto::new_with_type(InputType::Xtream))), BlockType::InputM3u => BlockInstance::Input(Rc::new(ConfigInputDto::new_with_type(InputType::M3u))), BlockType::InputLibrary => BlockInstance::Input(Rc::new(ConfigInputDto::new_with_type(InputType::Library))), + BlockType::InputEmby => BlockInstance::Input(Rc::new(ConfigInputDto::new_with_type(InputType::Emby))), + BlockType::InputJellyfin => BlockInstance::Input(Rc::new(ConfigInputDto::new_with_type(InputType::Jellyfin))), + BlockType::InputPlex => BlockInstance::Input(Rc::new(ConfigInputDto::new_with_type(InputType::Plex))), BlockType::Target => { let dto = ConfigTargetDto { name: String::new(), @@ -320,6 +323,15 @@ fn normalize_input_type_by_url(input: &mut ConfigInputDto, block_type: BlockType BlockType::InputLibrary => { input.input_type = InputType::Library; } + BlockType::InputEmby => { + input.input_type = InputType::Emby; + } + BlockType::InputJellyfin => { + input.input_type = InputType::Jellyfin; + } + BlockType::InputPlex => { + input.input_type = InputType::Plex; + } _ => {} } } @@ -2085,3 +2097,45 @@ fn update_pending_line(line: &Element, from: Position, to: Position) -> Result<( line.set_attribute("x2", &to.0.to_string())?; line.set_attribute("y2", &to.1.to_string()) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn media_server_blocks_create_matching_input_types() { + let cases = [ + (BlockType::InputEmby, InputType::Emby), + (BlockType::InputJellyfin, InputType::Jellyfin), + (BlockType::InputPlex, InputType::Plex), + ]; + + for (block_type, expected_input_type) in cases { + let BlockInstance::Input(input) = create_instance(block_type) else { + panic!("media_server block should create an input instance"); + }; + assert_eq!(input.input_type, expected_input_type); + } + } + + #[test] + fn media_server_block_normalization_preserves_input_type() { + let cases = [ + (BlockType::InputEmby, InputType::Emby), + (BlockType::InputJellyfin, InputType::Jellyfin), + (BlockType::InputPlex, InputType::Plex), + ]; + + for (block_type, expected_input_type) in cases { + let mut input = ConfigInputDto { + input_type: InputType::Xtream, + url: "batch://should-not-convert-media-server".to_string(), + ..ConfigInputDto::default() + }; + + normalize_input_type_by_url(&mut input, block_type); + + assert_eq!(input.input_type, expected_input_type); + } + } +} diff --git a/frontend/src/app/components/source_editor/output_form.rs b/frontend/src/app/components/source_editor/output_form.rs index 65c8a022f..89f64de0d 100644 --- a/frontend/src/app/components/source_editor/output_form.rs +++ b/frontend/src/app/components/source_editor/output_form.rs @@ -22,7 +22,13 @@ pub fn ConfigOutputView(props: &ConfigOutputViewProps) -> Html { match &*source_editor_ctx.edit_mode { EditMode::Active(block_instance) => match block_instance.block_type { - BlockType::InputXtream | BlockType::InputM3u | BlockType::InputLibrary | BlockType::Target => html! {}, + BlockType::InputXtream + | BlockType::InputM3u + | BlockType::InputLibrary + | BlockType::InputEmby + | BlockType::InputJellyfin + | BlockType::InputPlex + | BlockType::Target => html! {}, BlockType::OutputM3u => { let output = props.output.as_ref().and_then(|to| { if let TargetOutputDto::M3u(m3u) = &**to { diff --git a/frontend/src/app/components/source_editor/sidebar.rs b/frontend/src/app/components/source_editor/sidebar.rs index 30fb6f365..7b7f1708d 100644 --- a/frontend/src/app/components/source_editor/sidebar.rs +++ b/frontend/src/app/components/source_editor/sidebar.rs @@ -4,7 +4,15 @@ use crate::{ }; use yew::prelude::*; -pub const BLOCK_TYPES_INPUT: [BlockType; 3] = [BlockType::InputXtream, BlockType::InputM3u, BlockType::InputLibrary]; +pub const BLOCK_TYPES_INPUT: [BlockType; 3] = [ + BlockType::InputXtream, + BlockType::InputM3u, + BlockType::InputLibrary, + // Enable when media_server feature fully implemented + // BlockType::InputJellyfin, + // BlockType::InputEmby, + // BlockType::InputPlex, +]; pub const BLOCK_TYPES_TARGET: [BlockType; 1] = [BlockType::Target]; diff --git a/frontend/src/services/config_service.rs b/frontend/src/services/config_service.rs index 814ec20fb..b2582ec33 100644 --- a/frontend/src/services/config_service.rs +++ b/frontend/src/services/config_service.rs @@ -42,10 +42,7 @@ pub struct SetupWebUserCredentialDto { impl fmt::Debug for SetupWebUserCredentialDto { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - f.debug_struct("SetupWebUserCredentialDto") - .field("username", &self.username) - .field("password", &"") - .finish() + f.debug_struct("SetupWebUserCredentialDto").field("username", &self.username).field("password", &"***").finish() } } diff --git a/shared/src/model/config/input.rs b/shared/src/model/config/input.rs index 2a980efe6..54fe521e3 100644 --- a/shared/src/model/config/input.rs +++ b/shared/src/model/config/input.rs @@ -8,9 +8,10 @@ use crate::{ arc_str_serde, arc_str_vec_serde, default_as_true, default_probe_delay_secs, default_probe_live_interval, default_resolve_background, default_resolve_delay_secs, default_xtream_live_stream_use_prefix, deserialize_timestamp, get_credentials_from_url_str, get_trimmed_string, is_blank_optional_string, - is_default_probe_delay_secs, is_default_probe_live_interval, is_default_resolve_delay_secs, is_false, is_true, - is_zero_i16, is_zero_u16, parse_duration_seconds, parse_provider_scheme_url_parts, sanitize_sensitive_info, - serialize_option_vec_flow_map_items, trim_last_slash, Internable, BATCH_SCHEME_PREFIX, PROVIDER_SCHEME_PREFIX, + is_default_probe_delay_secs, is_default_probe_live_interval, is_default_resolve_delay_secs, is_false, + is_non_blank_optional_string, is_true, is_zero_i16, is_zero_u16, parse_duration_seconds, + parse_provider_scheme_url_parts, sanitize_sensitive_info, serialize_option_vec_flow_map_items, trim_last_slash, + Internable, BATCH_SCHEME_PREFIX, PROVIDER_SCHEME_PREFIX, }, }; use enum_iterator::Sequence; @@ -94,6 +95,12 @@ pub enum InputType { XtreamBatch, #[serde(rename = "library")] Library, + #[serde(rename = "emby")] + Emby, + #[serde(rename = "jellyfin")] + Jellyfin, + #[serde(rename = "plex")] + Plex, } impl InputType { @@ -102,10 +109,17 @@ impl InputType { const M3U_BATCH: &'static str = "m3u_batch"; const XTREAM_BATCH: &'static str = "xtream_batch"; const LIBRARY: &'static str = "library"; + const EMBY: &'static str = "emby"; + const JELLYFIN: &'static str = "jellyfin"; + const PLEX: &'static str = "plex"; pub fn is_xtream(&self) -> bool { matches!(self, Self::Xtream | Self::XtreamBatch) } pub fn is_m3u(&self) -> bool { matches!(self, Self::M3u | Self::M3uBatch) } + pub fn uses_standard_input_url(&self) -> bool { + matches!(self, Self::M3u | Self::Xtream | Self::M3uBatch | Self::XtreamBatch) + } pub fn is_library(&self) -> bool { matches!(self, Self::Library) } + pub fn is_media_server(&self) -> bool { matches!(self, Self::Emby | Self::Jellyfin | Self::Plex) } } impl Display for InputType { @@ -119,6 +133,9 @@ impl Display for InputType { Self::M3uBatch => Self::M3U_BATCH, Self::XtreamBatch => Self::XTREAM_BATCH, Self::Library => Self::LIBRARY, + Self::Emby => Self::EMBY, + Self::Jellyfin => Self::JELLYFIN, + Self::Plex => Self::PLEX, } ) } @@ -138,6 +155,12 @@ impl FromStr for InputType { Ok(Self::XtreamBatch) } else if s.eq(Self::LIBRARY) { Ok(Self::Library) + } else if s.eq(Self::EMBY) { + Ok(Self::Emby) + } else if s.eq(Self::JELLYFIN) { + Ok(Self::Jellyfin) + } else if s.eq(Self::PLEX) { + Ok(Self::Plex) } else { Err(TuliproxError::ConfigInput(format!("Unknown InputType: {}", s))) } @@ -313,6 +336,304 @@ impl ConfigInputOptionsDto { } } +pub const fn default_media_server_catalog_page_size() -> u16 { 100 } +pub const fn default_media_server_catalog_request_delay_ms() -> u64 { 250 } +pub const fn is_default_media_server_catalog_page_size(value: &u16) -> bool { + *value == default_media_server_catalog_page_size() +} +pub const fn is_default_media_server_catalog_request_delay_ms(value: &u64) -> bool { + *value == default_media_server_catalog_request_delay_ms() +} + +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq, Default)] +pub enum MediaServerCatalogRefreshModeDto { + #[serde(rename = "manual")] + #[default] + Manual, + #[serde(rename = "scheduled")] + Scheduled, +} + +pub fn is_default_media_server_catalog_refresh_mode(value: &MediaServerCatalogRefreshModeDto) -> bool { + *value == MediaServerCatalogRefreshModeDto::default() +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct MediaServerCatalogConfigDto { + #[serde(default, skip_serializing_if = "is_default_media_server_catalog_refresh_mode")] + pub refresh_mode: MediaServerCatalogRefreshModeDto, + #[serde(default, skip_serializing_if = "is_false")] + pub refresh_on_startup: bool, + #[serde( + default = "default_media_server_catalog_page_size", + skip_serializing_if = "is_default_media_server_catalog_page_size" + )] + pub page_size: u16, + #[serde( + default = "default_media_server_catalog_request_delay_ms", + skip_serializing_if = "is_default_media_server_catalog_request_delay_ms" + )] + pub request_delay_ms: u64, + #[serde(default = "default_as_true", skip_serializing_if = "is_true")] + pub include_media_sources: bool, + #[serde(default, alias = "include_file_paths", skip_serializing_if = "is_false")] + pub include_paths: bool, + #[serde(default, skip_serializing_if = "is_false")] + pub include_user_state: bool, +} + +impl Default for MediaServerCatalogConfigDto { + fn default() -> Self { + Self { + refresh_mode: MediaServerCatalogRefreshModeDto::default(), + refresh_on_startup: false, + page_size: default_media_server_catalog_page_size(), + request_delay_ms: default_media_server_catalog_request_delay_ms(), + include_media_sources: default_as_true(), + include_paths: false, + include_user_state: false, + } + } +} + +impl MediaServerCatalogConfigDto { + pub fn is_default(&self) -> bool { self == &Self::default() } + + pub fn prepare(&self, input_name: &Arc) -> Result<(), TuliproxError> { + if self.page_size == 0 { + return Err(TuliproxError::ConfigInput(format!( + "media server catalog page_size must be greater than zero (input: {input_name})" + ))); + } + Ok(()) + } +} + +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq, Default)] +pub enum MediaServerPlaybackInfoPolicyDto { + #[serde(rename = "on_demand")] + #[default] + OnDemand, + #[serde(rename = "disabled")] + Disabled, +} + +pub fn is_default_media_server_playback_info_policy(value: &MediaServerPlaybackInfoPolicyDto) -> bool { + *value == MediaServerPlaybackInfoPolicyDto::default() +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct MediaServerPlaybackConfigDto { + #[serde(default, skip_serializing_if = "is_default_media_server_playback_info_policy")] + pub playback_info_policy: MediaServerPlaybackInfoPolicyDto, + #[serde(default, skip_serializing_if = "is_false")] + pub preflight_streams: bool, + #[serde(default = "default_as_true", skip_serializing_if = "is_true")] + pub direct_play_only: bool, + #[serde(default, skip_serializing_if = "is_false")] + pub allow_transcode: bool, +} + +impl Default for MediaServerPlaybackConfigDto { + fn default() -> Self { + Self { + playback_info_policy: MediaServerPlaybackInfoPolicyDto::default(), + preflight_streams: false, + direct_play_only: default_as_true(), + allow_transcode: false, + } + } +} + +impl MediaServerPlaybackConfigDto { + pub fn is_default(&self) -> bool { self == &Self::default() } +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct MediaServerEnrichmentConfigDto { + #[serde(default, skip_serializing_if = "is_false")] + pub ffprobe: bool, + #[serde(default, skip_serializing_if = "is_false")] + pub tmdb_lookup: bool, + #[serde(default, skip_serializing_if = "is_false")] + pub fetch_images: bool, +} + +impl MediaServerEnrichmentConfigDto { + pub fn is_default(&self) -> bool { self == &Self::default() } +} + +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq, Default)] +pub enum MediaServerImagePolicyDto { + #[serde(rename = "proxy_on_demand")] + #[default] + ProxyOnDemand, + #[serde(rename = "disabled")] + Disabled, +} + +pub fn is_default_media_server_image_policy(value: &MediaServerImagePolicyDto) -> bool { + *value == MediaServerImagePolicyDto::default() +} + +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq)] +pub enum MediaServerLibraryKindDto { + #[serde(rename = "movies")] + Movies, + #[serde(rename = "tvshows", alias = "shows", alias = "series")] + TvShows, +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct MediaServerLibrarySelectorDetailsDto { + #[serde(default, skip_serializing_if = "is_blank_optional_string")] + pub id: Option, + #[serde(default, skip_serializing_if = "is_blank_optional_string")] + pub key: Option, + #[serde(default, skip_serializing_if = "is_blank_optional_string")] + pub name: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub kind: Option, +} + +impl MediaServerLibrarySelectorDetailsDto { + fn prepare(&mut self) { + self.id = get_trimmed_string(self.id.as_deref()); + self.key = get_trimmed_string(self.key.as_deref()); + self.name = get_trimmed_string(self.name.as_deref()); + } + + fn is_empty(&self) -> bool { + self.id.as_ref().is_none_or(|s| s.trim().is_empty()) + && self.key.as_ref().is_none_or(|s| s.trim().is_empty()) + && self.name.as_ref().is_none_or(|s| s.trim().is_empty()) + && self.kind.is_none() + } +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)] +#[serde(untagged)] +pub enum MediaServerLibrarySelectorDto { + Name(String), + Detailed(MediaServerLibrarySelectorDetailsDto), +} + +impl MediaServerLibrarySelectorDto { + fn prepare(&mut self) { + match self { + Self::Name(name) => *name = name.trim().to_string(), + Self::Detailed(details) => details.prepare(), + } + } + + pub fn is_empty(&self) -> bool { + match self { + Self::Name(name) => name.trim().is_empty(), + Self::Detailed(details) => details.is_empty(), + } + } +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct MediaServerInputConfigDto { + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub libraries: Vec, + #[serde(default, skip_serializing_if = "MediaServerCatalogConfigDto::is_default")] + pub catalog: MediaServerCatalogConfigDto, + #[serde(default, skip_serializing_if = "MediaServerPlaybackConfigDto::is_default")] + pub playback: MediaServerPlaybackConfigDto, + #[serde(default, skip_serializing_if = "MediaServerEnrichmentConfigDto::is_default")] + pub enrichment: MediaServerEnrichmentConfigDto, + #[serde(default, skip_serializing_if = "is_default_media_server_image_policy")] + pub image_policy: MediaServerImagePolicyDto, + #[serde(default, skip_serializing_if = "is_blank_optional_string")] + pub token: Option, + #[serde(default, skip_serializing_if = "is_blank_optional_string")] + pub api_key: Option, + #[serde(default, skip_serializing_if = "is_blank_optional_string")] + pub user_id: Option, + #[serde(default, skip_serializing_if = "is_blank_optional_string")] + pub account_token: Option, + #[serde(default, skip_serializing_if = "is_blank_optional_string")] + pub server_id: Option, + #[serde(default, skip_serializing_if = "is_blank_optional_string")] + pub machine_id: Option, + #[serde(default, skip_serializing_if = "is_blank_optional_string")] + pub server_name: Option, + #[serde(default = "default_as_true", skip_serializing_if = "is_true")] + pub prefer_https: bool, + #[serde(default, skip_serializing_if = "is_false")] + pub allow_relay: bool, +} + +impl Default for MediaServerInputConfigDto { + fn default() -> Self { + Self { + libraries: Vec::new(), + catalog: MediaServerCatalogConfigDto::default(), + playback: MediaServerPlaybackConfigDto::default(), + enrichment: MediaServerEnrichmentConfigDto::default(), + image_policy: MediaServerImagePolicyDto::default(), + token: None, + api_key: None, + user_id: None, + account_token: None, + server_id: None, + machine_id: None, + server_name: None, + prefer_https: default_as_true(), + allow_relay: false, + } + } +} + +impl MediaServerInputConfigDto { + pub fn normalize(&mut self) { + self.token = get_trimmed_string(self.token.as_deref()); + self.api_key = get_trimmed_string(self.api_key.as_deref()); + self.user_id = get_trimmed_string(self.user_id.as_deref()); + self.account_token = get_trimmed_string(self.account_token.as_deref()); + self.server_id = get_trimmed_string(self.server_id.as_deref()); + self.machine_id = get_trimmed_string(self.machine_id.as_deref()); + self.server_name = get_trimmed_string(self.server_name.as_deref()); + + for library in &mut self.libraries { + library.prepare(); + } + } + + pub fn prepare(&mut self, input_name: &Arc) -> Result<(), TuliproxError> { + self.normalize(); + self.catalog.prepare(input_name)?; + + if self.libraries.iter().any(MediaServerLibrarySelectorDto::is_empty) { + return Err(TuliproxError::ConfigInput(format!( + "media_server library selectors must not be empty (input: {input_name})" + ))); + } + Ok(()) + } + + pub fn has_any_emby_jellyfin_auth(&self) -> bool { + is_non_blank_optional_string(&self.token) || is_non_blank_optional_string(&self.api_key) + } + + pub fn has_any_plex_token(&self) -> bool { + is_non_blank_optional_string(&self.account_token) || is_non_blank_optional_string(&self.token) + } + + pub fn has_plex_server_selector(&self) -> bool { + is_non_blank_optional_string(&self.server_id) + || is_non_blank_optional_string(&self.machine_id) + || is_non_blank_optional_string(&self.server_name) + } +} + #[derive(Debug, Default, Copy, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)] pub enum ClusterSource { #[serde(rename = "staged")] @@ -495,6 +816,8 @@ pub struct ConfigInputDto { pub enabled: bool, #[serde(default, skip_serializing_if = "Option::is_none")] pub options: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub media_server: Option, #[serde(default, skip_serializing_if = "is_blank_optional_string")] pub cache_duration: Option, #[serde(skip)] @@ -531,6 +854,7 @@ impl Default for ConfigInputDto { persist: None, enabled: default_as_true(), options: None, + media_server: None, cache_duration: None, cache_duration_seconds: 0, aliases: None, @@ -566,9 +890,105 @@ impl ConfigInputDto { } } InputType::Library => InputType::Library, + InputType::Emby => InputType::Emby, + InputType::Jellyfin => InputType::Jellyfin, + InputType::Plex => InputType::Plex, }; } + fn prepare_media_server_input(&mut self) -> Result<(), TuliproxError> { + if !self.input_type.is_media_server() { + return Ok(()); + } + + let trimmed_url = self.url.trim(); + if trimmed_url.starts_with(BATCH_SCHEME_PREFIX) || trimmed_url.starts_with(PROVIDER_SCHEME_PREFIX) { + return Err(TuliproxError::ConfigInput(format!( + "media-server input does not support batch:// or provider:// URLs (input: {})", + self.name + ))); + } + if self.aliases.as_ref().is_some_and(|aliases| !aliases.is_empty()) { + return Err(TuliproxError::ConfigInput(format!( + "media-server input does not support aliases (input: {})", + self.name + ))); + } + if self.staged.as_ref().is_some_and(|staged| staged.enabled) { + return Err(TuliproxError::ConfigInput(format!( + "media-server input does not support staged inputs (input: {})", + self.name + ))); + } + if self.epg.is_some() { + return Err(TuliproxError::ConfigInput(format!( + "media-server input does not support EPG configuration (input: {})", + self.name + ))); + } + if self.panel_api.is_some() { + return Err(TuliproxError::ConfigInput(format!( + "media-server input does not support panel_api configuration (input: {})", + self.name + ))); + } + if self.provider.as_ref().is_some_and(|provider| !provider.is_empty()) { + return Err(TuliproxError::ConfigInput(format!( + "media-server input does not support provider failover definitions (input: {})", + self.name + ))); + } + let Some(media_server) = self.media_server.as_mut() else { + return Err(TuliproxError::ConfigInput(format!( + "media_server configuration is mandatory for input type {} (input: {})", + self.input_type, self.name + ))); + }; + media_server.prepare(&self.name)?; + if media_server.libraries.is_empty() { + return Err(TuliproxError::ConfigInput(format!( + "media-server input requires at least one selected library (input: {})", + self.name + ))); + } + + match self.input_type { + InputType::Emby | InputType::Jellyfin => { + if trimmed_url.is_empty() { + return Err(TuliproxError::ConfigInput(format!( + "url is mandatory for input type {} (input: {})", + self.input_type, self.name + ))); + } + let has_login = self.username.as_ref().is_some_and(|u| !u.trim().is_empty()) + && self.password.as_ref().is_some_and(|p| !p.trim().is_empty()); + if !media_server.has_any_emby_jellyfin_auth() && !has_login { + return Err(TuliproxError::ConfigInput(format!( + "media-server input type {} requires media_server token/api_key or username/password bootstrap credentials (input: {})", + self.input_type, self.name + ))); + } + } + InputType::Plex => { + if !media_server.has_any_plex_token() { + return Err(TuliproxError::ConfigInput(format!( + "media-server input type plex requires media_server.account_token or media_server.token (input: {})", + self.name + ))); + } + if trimmed_url.is_empty() && !media_server.has_plex_server_selector() { + return Err(TuliproxError::ConfigInput(format!( + "media-server input type plex requires a server selector such as media_server.machine_id, media_server.server_id, or media_server.server_name when input.url is omitted (input: {})", + self.name + ))); + } + } + InputType::M3u | InputType::Xtream | InputType::M3uBatch | InputType::XtreamBatch | InputType::Library => {} + } + + Ok(()) + } + #[allow(clippy::cast_possible_truncation)] pub fn prepare( &mut self, @@ -590,6 +1010,12 @@ impl ConfigInputDto { self.url = self.url.trim().to_string(); self.normalize_input_type_from_batch_url(); + if let Some(media_server) = self.media_server.as_mut() { + media_server.normalize(); + } + if self.enabled { + self.prepare_media_server_input()?; + } if self.url.starts_with(PROVIDER_SCHEME_PREFIX) && matches!(self.input_type, InputType::M3uBatch | InputType::XtreamBatch) { @@ -779,6 +1205,7 @@ impl ConfigInputDto { } pub fn prepare_type(&mut self) -> Result<(), TuliproxError> { + self.url = self.url.trim().to_string(); self.normalize_input_type_from_batch_url(); if self.url.starts_with(PROVIDER_SCHEME_PREFIX) && matches!(self.input_type, InputType::M3uBatch | InputType::XtreamBatch) @@ -1017,6 +1444,267 @@ mod tests { ConfigInputDto { name: "test_input".intern(), ..ConfigInputDto::default() } } + fn media_server_config_with_library() -> MediaServerInputConfigDto { + MediaServerInputConfigDto { + libraries: vec![MediaServerLibrarySelectorDto::Name("Movies".to_string())], + ..MediaServerInputConfigDto::default() + } + } + + #[test] + fn prepare_rejects_blank_media_server_credentials_and_selectors() { + let mut emby = ConfigInputDto { + name: "emby_media_server".intern(), + input_type: InputType::Emby, + url: "https://media.example.invalid".to_string(), + media_server: Some(MediaServerInputConfigDto { + token: Some(" ".to_string()), + api_key: Some("".to_string()), + ..media_server_config_with_library() + }), + ..ConfigInputDto::default() + }; + let err = prepare_dto(&mut emby).expect_err("blank token/api_key should be rejected"); + assert!(err.to_string().contains("requires media_server token/api_key")); + + let mut plex = ConfigInputDto { + name: "plex_media_server".intern(), + input_type: InputType::Plex, + media_server: Some(MediaServerInputConfigDto { + account_token: Some(" ".to_string()), + server_id: Some(" ".to_string()), + ..media_server_config_with_library() + }), + ..ConfigInputDto::default() + }; + let err = prepare_dto(&mut plex).expect_err("blank plex token should be rejected"); + assert!(err.to_string().contains("requires media_server.account_token or media_server.token")); + } + + #[test] + fn prepare_accepts_media_server_max_connections_as_stream_limit() { + let mut dto = ConfigInputDto { + name: "emby_media_server".intern(), + input_type: InputType::Emby, + url: "https://media.example.invalid".to_string(), + max_connections: 1, + media_server: Some(MediaServerInputConfigDto { + token: Some("token-value".to_string()), + ..media_server_config_with_library() + }), + ..ConfigInputDto::default() + }; + + prepare_dto(&mut dto).expect("media_server inputs reuse max_connections stream-limit semantics"); + } + + fn prepare_dto(dto: &mut ConfigInputDto) -> Result { + dto.prepare(0, false, &HashSet::new(), None) + } + + #[test] + fn media_server_defaults_are_conservative() { + let media_server = MediaServerInputConfigDto::default(); + + assert_eq!(media_server.catalog.page_size, 100); + assert_eq!(media_server.catalog.request_delay_ms, 250); + assert!(media_server.catalog.include_media_sources); + assert!(!media_server.catalog.include_paths); + assert!(!media_server.catalog.include_user_state); + assert!(!media_server.catalog.refresh_on_startup); + assert!(media_server.playback.direct_play_only); + assert!(!media_server.playback.allow_transcode); + assert!(!media_server.playback.preflight_streams); + assert!(!media_server.enrichment.ffprobe); + assert!(!media_server.enrichment.tmdb_lookup); + assert!(!media_server.enrichment.fetch_images); + assert_eq!(media_server.image_policy, MediaServerImagePolicyDto::ProxyOnDemand); + assert!(!media_server.allow_relay); + } + + #[test] + fn prepare_accepts_emby_media_server_with_token_and_library() { + let mut dto = ConfigInputDto { + name: "emby_media_server".intern(), + input_type: InputType::Emby, + url: " https://media.example.invalid/ ".to_string(), + media_server: Some(MediaServerInputConfigDto { + token: Some(" token-value ".to_string()), + ..media_server_config_with_library() + }), + ..ConfigInputDto::default() + }; + + prepare_dto(&mut dto).expect("emby media_server config should prepare"); + + assert_eq!(dto.url, "https://media.example.invalid/"); + assert!(dto.input_type.is_media_server()); + assert_eq!( + dto.media_server.as_ref().and_then(|media_server| media_server.token.as_deref()), + Some("token-value") + ); + } + + #[test] + fn prepare_rejects_media_server_without_media_server_block() { + let mut dto = ConfigInputDto { + name: "emby_media_server".intern(), + input_type: InputType::Emby, + url: "https://media.example.invalid".to_string(), + ..ConfigInputDto::default() + }; + + let err = prepare_dto(&mut dto).expect_err("media_server block should be mandatory"); + assert!(err.to_string().contains("media_server configuration is mandatory")); + } + + #[test] + fn prepare_allows_disabled_media_server_input_with_incomplete_config() { + let mut dto = ConfigInputDto { + name: "disabled_plex".intern(), + input_type: InputType::Plex, + enabled: false, + ..ConfigInputDto::default() + }; + + prepare_dto(&mut dto).expect("disabled media_server input should not require active playback/catalog config"); + assert_eq!(dto.input_type, InputType::Plex); + assert!(!dto.enabled); + } + + #[test] + fn prepare_normalizes_disabled_media_server_config_without_enforcing_invariants() { + let mut dto = ConfigInputDto { + name: "disabled_emby".intern(), + input_type: InputType::Emby, + enabled: false, + media_server: Some(MediaServerInputConfigDto { + token: Some(" token-value ".to_string()), + libraries: vec![MediaServerLibrarySelectorDto::Name(" ".to_string())], + catalog: MediaServerCatalogConfigDto { page_size: 0, ..MediaServerCatalogConfigDto::default() }, + ..MediaServerInputConfigDto::default() + }), + ..ConfigInputDto::default() + }; + + prepare_dto(&mut dto).expect("disabled media_server input can preserve incomplete config for later repair"); + let media_server = dto.media_server.as_ref().expect("media_server config should be preserved"); + assert_eq!(media_server.token.as_deref(), Some("token-value")); + assert!(media_server.libraries[0].is_empty()); + } + + #[test] + fn prepare_rejects_emby_media_server_without_input_url() { + let mut dto = ConfigInputDto { + name: "emby_media_server".intern(), + input_type: InputType::Emby, + media_server: Some(MediaServerInputConfigDto { + token: Some("token-value".to_string()), + ..media_server_config_with_library() + }), + ..ConfigInputDto::default() + }; + + let err = prepare_dto(&mut dto).expect_err("emby media_server input should require a direct server URL"); + assert!(err.to_string().contains("url is mandatory for input type emby")); + } + + #[test] + fn prepare_rejects_media_server_provider_scheme_url() { + let mut dto = ConfigInputDto { + name: "emby_media_server".intern(), + input_type: InputType::Emby, + url: " provider://media-server ".to_string(), + media_server: Some(MediaServerInputConfigDto { + token: Some("token-value".to_string()), + ..media_server_config_with_library() + }), + ..ConfigInputDto::default() + }; + + let err = prepare_dto(&mut dto).expect_err("media_server input must not use provider URLs"); + assert!(err.to_string().contains("does not support batch:// or provider://")); + } + + #[test] + fn prepare_rejects_media_server_staged_input() { + let mut dto = ConfigInputDto { + name: "emby_media_server".intern(), + input_type: InputType::Emby, + url: "https://media.example.invalid".to_string(), + media_server: Some(MediaServerInputConfigDto { + token: Some("token-value".to_string()), + ..media_server_config_with_library() + }), + staged: Some(StagedInputDto { enabled: true, name: "staged".intern(), ..StagedInputDto::default() }), + ..ConfigInputDto::default() + }; + + let err = prepare_dto(&mut dto).expect_err("media_server input must reject staged config"); + assert!(err.to_string().contains("does not support staged inputs")); + } + + #[test] + fn prepare_rejects_plex_without_token_or_server_selector() { + let mut without_token = ConfigInputDto { + name: "plex_media_server".intern(), + input_type: InputType::Plex, + media_server: Some(MediaServerInputConfigDto { + machine_id: Some("machine".to_string()), + ..media_server_config_with_library() + }), + ..ConfigInputDto::default() + }; + let err = prepare_dto(&mut without_token).expect_err("plex token should be mandatory"); + assert!(err.to_string().contains("requires media_server.account_token or media_server.token")); + + let mut without_selector = ConfigInputDto { + name: "plex_media_server".intern(), + input_type: InputType::Plex, + media_server: Some(MediaServerInputConfigDto { + account_token: Some("token".to_string()), + ..media_server_config_with_library() + }), + ..ConfigInputDto::default() + }; + let err = prepare_dto(&mut without_selector).expect_err("plex server selector should be mandatory"); + assert!(err.to_string().contains("requires a server selector")); + } + + #[test] + fn prepare_accepts_plex_without_input_url_when_discovery_is_configured() { + let mut dto = ConfigInputDto { + name: "plex_media_server".intern(), + input_type: InputType::Plex, + media_server: Some(MediaServerInputConfigDto { + account_token: Some("token".to_string()), + machine_id: Some("machine".to_string()), + ..media_server_config_with_library() + }), + ..ConfigInputDto::default() + }; + + prepare_dto(&mut dto).expect("plex discovery config should not require input.url"); + assert_eq!(dto.input_type, InputType::Plex); + } + + #[test] + fn prepare_accepts_plex_media_server_with_direct_url_without_selector() { + let mut dto = ConfigInputDto { + name: "plex_media_server".intern(), + input_type: InputType::Plex, + url: "https://plex.example.invalid".to_string(), + media_server: Some(MediaServerInputConfigDto { + token: Some("token".to_string()), + ..media_server_config_with_library() + }), + ..ConfigInputDto::default() + }; + + prepare_dto(&mut dto).expect("direct Plex URL should not require MyPlex server selector"); + assert_eq!(dto.input_type, InputType::Plex); + } + #[test] fn test_epg_url_from_explicit_main_credentials() { let mut dto = create_test_dto(); @@ -1264,6 +1952,20 @@ mod tests { assert_eq!(dto.input_type, InputType::Xtream); } + #[test] + fn prepare_type_does_not_validate_media_server_config() { + let mut dto = ConfigInputDto { + name: "emby_media_server".intern(), + input_type: InputType::Emby, + url: "https://media.example.invalid".to_string(), + ..ConfigInputDto::default() + }; + + dto.prepare_type().expect("prepare_type only normalizes type/url"); + let err = prepare_dto(&mut dto).expect_err("full prepare should validate missing media_server block"); + assert!(err.to_string().contains("media_server configuration is mandatory")); + } + #[test] fn prepare_batch_url_does_not_require_xtream_credentials() { let mut dto = ConfigInputDto { diff --git a/shared/src/model/config/macros.rs b/shared/src/model/config/macros.rs index 2a635bbc3..8e9b19644 100644 --- a/shared/src/model/config/macros.rs +++ b/shared/src/model/config/macros.rs @@ -1,15 +1,17 @@ #[macro_export] macro_rules! check_input_credentials { ($this:ident, $input_type:expr, $definition:expr, $alias:expr ) => { - let __tp_input_name = $this.name.to_string(); - let __tp_input_name = __tp_input_name.trim().to_string(); - let __tp_input_name_suffix = - if __tp_input_name.is_empty() { String::new() } else { format!(" (input: {})", __tp_input_name) }; + let input_name = $this.name.to_string(); + let input_name = input_name.trim().to_string(); + let input_name_suffix = + if input_name.is_empty() { String::new() } else { format!(" (input: {input_name})") }; if !matches!($input_type, InputType::Library) { $this.url = $this.url.trim().to_string(); - if $this.url.is_empty() { - return Err($crate::error::TuliproxError::ConfigInput(format!("url for input is mandatory{}", __tp_input_name_suffix))); + // This generic check only applies to classic URL-backed playlist inputs. Media-server + // inputs have provider-specific URL/discovery validation in their prepare path. + if $input_type.uses_standard_input_url() && $this.url.is_empty() { + return Err($crate::error::TuliproxError::ConfigInput(format!("url for input is mandatory{input_name_suffix}"))); } $this.username = $crate::utils::get_trimmed_string($this.username.as_deref()); @@ -28,7 +30,7 @@ macro_rules! check_input_credentials { InputType::M3uBatch => { if $definition { if $this.url.trim().is_empty() { - return Err($crate::error::TuliproxError::ConfigInput(format!("for input type m3u-batch: url is mandatory{}", __tp_input_name_suffix))); + return Err($crate::error::TuliproxError::ConfigInput(format!("for input type m3u-batch: url is mandatory{input_name_suffix}"))); } } @@ -40,8 +42,7 @@ macro_rules! check_input_credentials { InputType::Xtream => { if $this.username.is_none() || $this.password.is_none() { return Err($crate::error::TuliproxError::ConfigInput(format!( - "for input type xtream: username and password are mandatory{}", - __tp_input_name_suffix + "for input type xtream: username and password are mandatory{input_name_suffix}", ))); } } @@ -49,8 +50,7 @@ macro_rules! check_input_credentials { if $definition { if $this.url.trim().is_empty() { return Err($crate::error::TuliproxError::ConfigInput(format!( - "for input type xtream-batch: url is mandatory{}", - __tp_input_name_suffix + "for input type xtream-batch: url is mandatory{input_name_suffix}", ))); } } @@ -64,20 +64,19 @@ macro_rules! check_input_credentials { if is_batch_url { if has_credentials { return Err($crate::error::TuliproxError::ConfigInput(format!( - "input type xtream-batch with batch:// URL should not define username or password attribute{}", - __tp_input_name_suffix + "input type xtream-batch with batch:// URL should not define username or password attribute{input_name_suffix}", ))); } } else if !has_username || !has_password { return Err($crate::error::TuliproxError::ConfigInput(format!( - "for input type xtream-batch without batch:// URL: username and password are mandatory{}", - __tp_input_name_suffix + "for input type xtream-batch without batch:// URL: username and password are mandatory{input_name_suffix}", ))); } } } - InputType::Library => { - // nothing to do + InputType::Library | InputType::Emby | InputType::Jellyfin | InputType::Plex => { + // Media-server credentials live in the dedicated media_server block; detailed + // validation happens in ConfigInputDto/ConfigInput prepare methods. } } }; @@ -86,10 +85,9 @@ macro_rules! check_input_credentials { #[macro_export] macro_rules! check_input_connections { ($this:ident, $input_type:expr, $alias:expr) => { - let __tp_input_name = $this.name.to_string(); - let __tp_input_name = __tp_input_name.trim().to_string(); - let __tp_input_name_suffix = - if __tp_input_name.is_empty() { String::new() } else { format!(" (input: {})", __tp_input_name) }; + let input_name = $this.name.to_string(); + let input_name = input_name.trim().to_string(); + let input_name_suffix = if input_name.is_empty() { String::new() } else { format!(" (input: {input_name})") }; match $input_type { InputType::M3u | InputType::Xtream => {} @@ -97,14 +95,12 @@ macro_rules! check_input_connections { if !$alias { if $this.max_connections > 0 { return Err($crate::error::TuliproxError::ConfigInput(format!( - "input type m3u-batch should not define max_connections attribute{}", - __tp_input_name_suffix + "input type m3u-batch should not define max_connections attribute{input_name_suffix}", ))); } if $this.priority != 0 { return Err($crate::error::TuliproxError::ConfigInput(format!( - "input type m3u-batch should not define priority attribute{}", - __tp_input_name_suffix + "input type m3u-batch should not define priority attribute{input_name_suffix}", ))); } } @@ -113,19 +109,17 @@ macro_rules! check_input_connections { if !$alias { if $this.max_connections > 0 { return Err($crate::error::TuliproxError::ConfigInput(format!( - "input type xtream-batch should not define max_connections attribute{}", - __tp_input_name_suffix + "input type xtream-batch should not define max_connections attribute{input_name_suffix}", ))); } if $this.priority != 0 { return Err($crate::error::TuliproxError::ConfigInput(format!( - "input type xtream-batch should not define priority attribute{}", - __tp_input_name_suffix + "input type xtream-batch should not define priority attribute{input_name_suffix}", ))); } } } - InputType::Library => {} + InputType::Library | InputType::Emby | InputType::Jellyfin | InputType::Plex => {} } }; } diff --git a/shared/src/model/xtream.rs b/shared/src/model/xtream.rs index 8f6bf757a..394a93249 100644 --- a/shared/src/model/xtream.rs +++ b/shared/src/model/xtream.rs @@ -29,8 +29,8 @@ impl fmt::Debug for XtreamLoginRequest { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { f.debug_struct("XtreamLoginRequest") .field("url", &self.url) - .field("username", &"") - .field("password", &"") + .field("username", &"***") + .field("password", &"***") .finish() } } diff --git a/shared/src/utils/default_utils.rs b/shared/src/utils/default_utils.rs index 941e06ad0..cdece3885 100644 --- a/shared/src/utils/default_utils.rs +++ b/shared/src/utils/default_utils.rs @@ -22,6 +22,8 @@ pub fn is_blank_optional_string(s: &Option) -> bool { s.as_ref().is_none_or(|s| s.chars().all(|c| c.is_whitespace())) } +pub fn is_non_blank_optional_string(s: &Option) -> bool { !is_blank_optional_string(s) } + pub fn is_blank_optional_arc_str(s: &Option>) -> bool { s.as_ref().is_none_or(|s| s.chars().all(|c| c.is_whitespace())) }