feat: add media-server input foundations (#742)

* feat: add remote media input foundations
* fixed macro variable names
* Added todo comment for source editor sidebar
This commit is contained in:
Juanjo Presa
2026-05-04 20:22:51 +02:00
committed by GitHub
parent ef8350a68c
commit 8a2e379f6d
39 changed files with 3376 additions and 56 deletions
@@ -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<String> = errors.iter().map(ToString::to_string).collect();
@@ -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 {
@@ -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,
+1 -1
View File
@@ -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", &"<redacted>")
.field("password", &"***")
.finish()
}
}
+307
View File
@@ -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<usize>,
pub fetched: usize,
}
impl MediaServerCatalogCursor {
pub fn from_page<T>(library: &MediaServerLibrary, page: &MediaServerPage<T>) -> 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<MediaServerLibrary>,
pub movies: Vec<MediaServerMovie>,
pub episodes: Vec<MediaServerEpisode>,
pub unsupported_libraries: Vec<MediaServerLibrary>,
}
impl MediaServerCatalogSnapshot {
pub fn item_count(&self) -> usize { self.movies.len() + self.episodes.len() }
}
#[derive(Debug, Clone, Default)]
pub struct MediaServerCatalogCache {
trusted: Option<MediaServerCatalogSnapshot>,
}
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<C>(
&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<C>(
client: &C,
policy: MediaServerCatalogRefreshPolicy,
) -> Result<MediaServerCatalogSnapshot, MediaServerError>
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<T>(library: &MediaServerLibrary, page: &MediaServerPage<T>) -> 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<Vec<Result<MediaServerPage<MediaServerMovie>, MediaServerError>>>,
episode_pages: Mutex<Vec<Result<MediaServerPage<MediaServerEpisode>, MediaServerError>>>,
libraries: Vec<MediaServerLibrary>,
}
impl MockMediaServerCatalogClient {
fn with_libraries(libraries: Vec<MediaServerLibrary>) -> Self {
Self { libraries, ..Self::default() }
}
}
impl MediaServerCatalogClient for MockMediaServerCatalogClient {
async fn discover(&self) -> Result<MediaServerStatus, MediaServerError> {
Ok(MediaServerStatus {
kind: MediaServerKind::Emby,
server_id: "server-redacted".into(),
display_name: None,
version: None,
owned: None,
})
}
async fn list_libraries(&self) -> Result<Vec<MediaServerLibrary>, MediaServerError> { Ok(self.libraries.clone()) }
async fn list_movies(
&self,
_library: &MediaServerLibraryRef,
_page: MediaServerPageRequest,
) -> Result<MediaServerPage<MediaServerMovie>, MediaServerError> {
self.movie_pages.lock().expect("lock").remove(0)
}
async fn list_episodes(
&self,
_library: &MediaServerLibraryRef,
_page: MediaServerPageRequest,
) -> Result<MediaServerPage<MediaServerEpisode>, MediaServerError> {
self.episode_pages.lock().expect("lock").remove(0)
}
async fn open_stream(
&self,
_stream_ref: &MediaServerStreamRef,
_range: Option<&str>,
) -> Result<crate::media_server::MediaServerStreamResponse, MediaServerError> {
Ok(empty_stream_response())
}
async fn open_image(
&self,
_image_ref: &MediaServerImageRef,
) -> Result<MediaServerResourceResponse, MediaServerError> {
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, MediaServerError>(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::<str>::from(id),
title: Arc::<str>::from("Movie Redacted"),
year: None,
source_version_hint: None,
provider_hints: Vec::<MediaServerProviderIdHint>::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);
}
}
+128
View File
@@ -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<MediaServerStatus, MediaServerError>;
async fn list_libraries(&self) -> Result<Vec<MediaServerLibrary>, MediaServerError>;
async fn list_movies(
&self,
library: &MediaServerLibraryRef,
page: MediaServerPageRequest,
) -> Result<MediaServerPage<MediaServerMovie>, MediaServerError>;
async fn list_episodes(
&self,
library: &MediaServerLibraryRef,
page: MediaServerPageRequest,
) -> Result<MediaServerPage<MediaServerEpisode>, MediaServerError>;
async fn open_stream(
&self,
stream_ref: &MediaServerStreamRef,
range: Option<&str>,
) -> Result<MediaServerStreamResponse, MediaServerError>;
async fn open_image(&self, image_ref: &MediaServerImageRef) -> Result<MediaServerResourceResponse, MediaServerError>;
}
#[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<reqwest::Response, MediaServerError> {
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);
}
}
+88
View File
@@ -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<String>,
pub server_name: Option<String>,
pub version: Option<String>,
}
#[derive(Debug, Clone, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "PascalCase", bound(deserialize = "T: Deserialize<'de>"))]
pub struct EmbyItemsPageDto<T> {
#[serde(default)]
pub items: Vec<T>,
#[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<String>,
pub collection_type: Option<String>,
pub type_: Option<String>,
}
#[derive(Debug, Clone, Deserialize, PartialEq)]
#[serde(rename_all = "PascalCase")]
pub struct EmbyItemDto {
pub id: String,
pub name: Option<String>,
pub type_: Option<String>,
pub production_year: Option<u32>,
pub parent_id: Option<String>,
pub series_id: Option<String>,
pub series_name: Option<String>,
pub parent_index_number: Option<u32>,
pub index_number: Option<u32>,
#[serde(default)]
pub provider_ids: HashMap<String, String>,
#[serde(default)]
pub image_tags: HashMap<String, String>,
#[serde(default)]
pub media_sources: Vec<EmbyMediaSourceDto>,
// Parsed only so the boundary can explicitly ignore it by default.
pub path: Option<String>,
pub user_data: Option<serde_json::Value>,
}
#[derive(Debug, Clone, Deserialize, PartialEq)]
#[serde(rename_all = "PascalCase")]
pub struct EmbyMediaSourceDto {
pub id: Option<String>,
pub container: Option<String>,
pub path: Option<String>,
pub supports_direct_play: Option<bool>,
pub supports_direct_stream: Option<bool>,
}
#[derive(Debug, Clone, Deserialize, PartialEq)]
#[serde(rename_all = "PascalCase")]
pub struct EmbyPlaybackInfoDto {
#[serde(default)]
pub media_sources: Vec<EmbyMediaSourceDto>,
pub play_session_id: Option<String>,
}
#[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<EmbyItemDto> = 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());
}
}
+1
View File
@@ -0,0 +1 @@
pub mod dto;
+190
View File
@@ -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<StatusCode>,
detail: Option<String>,
}
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<str>) -> 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("***"));
}
}
+26
View File
@@ -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<T> = crate::media_server::emby::dto::EmbyItemsPageDto<T>;
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<JellyfinViewDto> = 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"));
}
}
+1
View File
@@ -0,0 +1 @@
pub mod dto;
+27
View File
@@ -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::*;
+335
View File
@@ -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", &"<stream>")
.finish()
}
}
pub fn classify_playback_origin(
input_type: InputType,
item_type: PlaylistItemType,
input_name: &Arc<str>,
item_url: &str,
) -> Result<PlaybackOrigin, MediaServerError> {
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<C>(
client: &C,
stream_ref: &MediaServerStreamRef,
range: Option<&str>,
) -> Result<MediaServerProxyResponse, MediaServerError>
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<C>(
client: &C,
image_ref: &MediaServerImageRef,
) -> Result<MediaServerProxyResponse, MediaServerError>
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::<Bytes, MediaServerError>(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<str>, item_url: &str) -> Result<MediaServerStreamRef, MediaServerError> {
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<String> = 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::<str>::from(parts[1].as_str()),
item_id: Arc::<str>::from(parts[2].as_str()),
media_source_id: query_value(query, "media_source_id").map(Arc::<str>::from),
}),
"jellyfin" => Ok(MediaServerStreamRef::Jellyfin {
input_name: input_name.clone(),
server_id: Arc::<str>::from(parts[1].as_str()),
item_id: Arc::<str>::from(parts[2].as_str()),
media_source_id: query_value(query, "media_source_id").map(Arc::<str>::from),
}),
"plex" => Ok(MediaServerStreamRef::Plex {
input_name: input_name.clone(),
server_id: Arc::<str>::from(parts[1].as_str()),
rating_key: Arc::<str>::from(parts[2].as_str()),
part_key: query_value(query, "part_key")
.map(Arc::<str>::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<String> {
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<u8> {
Some(hex_value(high)? << 4 | hex_value(low)?)
}
fn hex_value(value: u8) -> Option<u8> {
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<Option<String>>,
stream_error: Option<MediaServerError>,
}
impl MediaServerCatalogClient for MockPlaybackClient {
async fn discover(&self) -> Result<MediaServerStatus, MediaServerError> { unreachable!() }
async fn list_libraries(&self) -> Result<Vec<MediaServerLibrary>, MediaServerError> { unreachable!() }
async fn list_movies(
&self,
_library: &MediaServerLibraryRef,
_page: MediaServerPageRequest,
) -> Result<MediaServerPage<MediaServerMovie>, MediaServerError> {
unreachable!()
}
async fn list_episodes(
&self,
_library: &MediaServerLibraryRef,
_page: MediaServerPageRequest,
) -> Result<MediaServerPage<MediaServerEpisode>, MediaServerError> {
unreachable!()
}
async fn open_stream(
&self,
_stream_ref: &MediaServerStreamRef,
range: Option<&str>,
) -> Result<crate::media_server::MediaServerStreamResponse, MediaServerError> {
*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, MediaServerError>(Bytes::from_static(b"data")) }).boxed(),
})
}
async fn open_image(&self, _image_ref: &MediaServerImageRef) -> Result<MediaServerResourceResponse, MediaServerError> {
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::<Vec<_>>().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::<str>::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);
}
}
+249
View File
@@ -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<PlaylistGroup> {
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<str>, library_id: &Arc<str>, item_id: &Arc<str>, 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::<MediaServerProviderIdHint>::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::<MediaServerProviderIdHint>::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"));
}
}
+198
View File
@@ -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<PlexResourceDto>,
}
#[derive(Debug, Clone, Deserialize, PartialEq)]
pub struct PlexResourceDto {
#[serde(rename = "@name")]
pub name: Option<String>,
#[serde(rename = "@product")]
pub product: Option<String>,
#[serde(rename = "@productVersion")]
pub product_version: Option<String>,
#[serde(rename = "@clientIdentifier")]
pub client_identifier: Option<String>,
#[serde(rename = "@machineIdentifier")]
pub machine_identifier: Option<String>,
#[serde(rename = "@owned", default)]
pub owned: Option<u8>,
#[serde(rename = "@accessToken")]
pub access_token: Option<String>,
#[serde(rename = "Connection", default)]
pub connections: Vec<PlexConnectionDto>,
}
#[derive(Debug, Clone, Deserialize, PartialEq)]
pub struct PlexConnectionDto {
#[serde(rename = "@protocol")]
pub protocol: Option<String>,
#[serde(rename = "@uri")]
pub uri: Option<String>,
#[serde(rename = "@local", default)]
pub local: Option<u8>,
#[serde(rename = "@relay", default)]
pub relay: Option<u8>,
}
#[derive(Debug, Clone, Deserialize, PartialEq)]
#[serde(rename = "MediaContainer")]
pub struct PlexSectionsDto {
#[serde(rename = "Directory", default)]
pub directories: Vec<PlexSectionDto>,
}
#[derive(Debug, Clone, Deserialize, PartialEq)]
pub struct PlexSectionDto {
#[serde(rename = "@key")]
pub key: Option<String>,
#[serde(rename = "@title")]
pub title: Option<String>,
#[serde(rename = "@type")]
pub section_type: Option<String>,
}
#[derive(Debug, Clone, Deserialize, PartialEq)]
#[serde(rename = "MediaContainer")]
pub struct PlexMediaContainerDto {
#[serde(rename = "@size")]
pub size: Option<usize>,
#[serde(rename = "@totalSize")]
pub total_size: Option<usize>,
#[serde(rename = "Video", default)]
pub videos: Vec<PlexVideoDto>,
#[serde(rename = "Directory", default)]
pub directories: Vec<PlexDirectoryDto>,
}
#[derive(Debug, Clone, Deserialize, PartialEq)]
pub struct PlexDirectoryDto {
#[serde(rename = "@ratingKey")]
pub rating_key: Option<String>,
#[serde(rename = "@key")]
pub key: Option<String>,
#[serde(rename = "@type")]
pub item_type: Option<String>,
#[serde(rename = "@title")]
pub title: Option<String>,
#[serde(rename = "@year")]
pub year: Option<u32>,
#[serde(rename = "Guid", default)]
pub guids: Vec<PlexGuidDto>,
}
#[derive(Debug, Clone, Deserialize, PartialEq)]
pub struct PlexVideoDto {
#[serde(rename = "@ratingKey")]
pub rating_key: Option<String>,
#[serde(rename = "@key")]
pub key: Option<String>,
#[serde(rename = "@type")]
pub item_type: Option<String>,
#[serde(rename = "@title")]
pub title: Option<String>,
#[serde(rename = "@year")]
pub year: Option<u32>,
#[serde(rename = "@guid")]
pub guid: Option<String>,
#[serde(rename = "@thumb")]
pub thumb: Option<String>,
#[serde(rename = "@art")]
pub art: Option<String>,
#[serde(rename = "@parentRatingKey")]
pub parent_rating_key: Option<String>,
#[serde(rename = "@grandparentRatingKey")]
pub grandparent_rating_key: Option<String>,
#[serde(rename = "@grandparentTitle")]
pub grandparent_title: Option<String>,
#[serde(rename = "@parentIndex")]
pub parent_index: Option<u32>,
#[serde(rename = "@index")]
pub index: Option<u32>,
#[serde(rename = "@addedAt")]
pub added_at: Option<i64>,
#[serde(rename = "@updatedAt")]
pub updated_at: Option<i64>,
#[serde(rename = "Guid", default)]
pub guids: Vec<PlexGuidDto>,
#[serde(rename = "Media", default)]
pub media: Vec<PlexMediaDto>,
}
#[derive(Debug, Clone, Deserialize, PartialEq, Eq)]
pub struct PlexGuidDto {
#[serde(rename = "@id")]
pub id: Option<String>,
}
#[derive(Debug, Clone, Deserialize, PartialEq)]
pub struct PlexMediaDto {
#[serde(rename = "@id")]
pub id: Option<String>,
#[serde(rename = "@container")]
pub container: Option<String>,
#[serde(rename = "@duration")]
pub duration: Option<u64>,
#[serde(rename = "@bitrate")]
pub bitrate: Option<u32>,
#[serde(rename = "@width")]
pub width: Option<u32>,
#[serde(rename = "@height")]
pub height: Option<u32>,
#[serde(rename = "Part", default)]
pub parts: Vec<PlexPartDto>,
}
#[derive(Debug, Clone, Deserialize, PartialEq, Eq)]
pub struct PlexPartDto {
#[serde(rename = "@id")]
pub id: Option<String>,
#[serde(rename = "@key")]
pub key: Option<String>,
#[serde(rename = "@size")]
pub size: Option<u64>,
#[serde(rename = "@file")]
pub file: Option<String>,
#[serde(rename = "@container")]
pub container: Option<String>,
}
#[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);
}
}
+1
View File
@@ -0,0 +1 @@
pub mod dto;
+170
View File
@@ -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(|_| "<non-utf8>".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"));
}
}
+65
View File
@@ -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#"
<MediaContainer size="1">
<Device name="Server Redacted" product="Plex Media Server" productVersion="1.0.0" clientIdentifier="client-redacted" machineIdentifier="machine-redacted" owned="0" accessToken="resource-token-redacted">
<Connection protocol="https" uri="https://pms.example.invalid" local="0" relay="0" />
</Device>
</MediaContainer>
"#;
pub const PLEX_SECTIONS_XML: &str = r#"
<MediaContainer size="3">
<Directory key="1" title="Movies" type="movie" />
<Directory key="2" title="Shows" type="show" />
<Directory key="3" title="Music" type="artist" />
</MediaContainer>
"#;
pub const PLEX_MOVIES_XML: &str = r#"
<MediaContainer size="1" totalSize="1">
<Video ratingKey="rating-redacted-1" key="/library/metadata/rating-redacted-1" type="movie" title="Movie Redacted" year="2024" thumb="/library/metadata/rating-redacted-1/thumb" addedAt="1700000000" updatedAt="1700000001">
<Guid id="tmdb://12345" />
<Media id="media-redacted-1" container="mkv" duration="7200000" bitrate="8000" width="1920" height="1080">
<Part id="part-redacted" key="/library/parts/part-redacted/file.mkv" file="/redacted/upstream/path/movie.mkv" size="1024" container="mkv" />
</Media>
</Video>
</MediaContainer>
"#;
+269
View File
@@ -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<InputType> for MediaServerKind {
type Error = &'static str;
fn try_from(value: InputType) -> Result<Self, Self::Error> {
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<str>,
pub display_name: Option<Arc<str>>,
pub version: Option<Arc<str>>,
pub owned: Option<bool>,
}
#[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<str>,
pub server_id: Arc<str>,
pub library_id: Arc<str>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MediaServerLibrary {
pub reference: MediaServerLibraryRef,
pub name: Arc<str>,
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<T> {
pub request: MediaServerPageRequest,
pub total: Option<usize>,
pub upstream_item_count: usize,
pub items: Vec<T>,
}
impl<T> MediaServerPage<T> {
pub fn new(request: MediaServerPageRequest, total: Option<usize>, items: Vec<T>) -> Self {
let upstream_item_count = items.len();
Self { request, total, upstream_item_count, items }
}
pub fn with_upstream_item_count(
request: MediaServerPageRequest,
total: Option<usize>,
upstream_item_count: usize,
items: Vec<T>,
) -> 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<MediaServerPageRequest> {
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<str>,
pub value: Arc<str>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MediaServerMovie {
pub input_name: Arc<str>,
pub server_id: Arc<str>,
pub library_id: Arc<str>,
pub item_id: Arc<str>,
pub title: Arc<str>,
pub year: Option<u32>,
pub source_version_hint: Option<Arc<str>>,
pub provider_hints: Vec<MediaServerProviderIdHint>,
pub stream_ref: Option<MediaServerStreamRef>,
pub image_ref: Option<MediaServerImageRef>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MediaServerEpisode {
pub input_name: Arc<str>,
pub server_id: Arc<str>,
pub library_id: Arc<str>,
pub item_id: Arc<str>,
pub series_id: Option<Arc<str>>,
pub series_title: Option<Arc<str>>,
pub title: Arc<str>,
pub season: Option<u32>,
pub episode: Option<u32>,
pub source_version_hint: Option<Arc<str>>,
pub provider_hints: Vec<MediaServerProviderIdHint>,
pub stream_ref: Option<MediaServerStreamRef>,
pub image_ref: Option<MediaServerImageRef>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum MediaServerStreamRef {
Emby {
input_name: Arc<str>,
server_id: Arc<str>,
item_id: Arc<str>,
media_source_id: Option<Arc<str>>,
},
Jellyfin {
input_name: Arc<str>,
server_id: Arc<str>,
item_id: Arc<str>,
media_source_id: Option<Arc<str>>,
},
Plex {
input_name: Arc<str>,
server_id: Arc<str>,
rating_key: Arc<str>,
part_key: Arc<str>,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum MediaServerImageRef {
Emby {
input_name: Arc<str>,
server_id: Arc<str>,
item_id: Arc<str>,
image_kind: Arc<str>,
tag: Option<Arc<str>>,
},
Jellyfin {
input_name: Arc<str>,
server_id: Arc<str>,
item_id: Arc<str>,
image_kind: Arc<str>,
tag: Option<Arc<str>>,
},
Plex {
input_name: Arc<str>,
server_id: Arc<str>,
rating_key: Arc<str>,
image_path: Arc<str>,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MediaServerPlaybackLease {
pub provider_kind: MediaServerKind,
pub lease_id: Arc<str>,
}
#[derive(Debug, Clone)]
pub struct MediaServerResourceResponse {
pub status: StatusCode,
pub headers: HeaderMap,
pub body: Bytes,
}
pub type BoxedMediaServerStream = BoxStream<'static, Result<Bytes, MediaServerError>>;
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", &"<stream>")
.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::<u8>::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)));
}
}
+354 -2
View File
@@ -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<ConfigInputOptions> =
LazyLock::new(|| ConfigInputOptions::from(&ConfigInputOptionsDto::default()));
#[derive(Debug, Clone)]
pub struct MediaServerInputConfig {
pub libraries: Vec<MediaServerLibrarySelectorDto>,
pub catalog: MediaServerCatalogConfigDto,
pub playback: MediaServerPlaybackConfigDto,
pub enrichment: MediaServerEnrichmentConfigDto,
pub image_policy: MediaServerImagePolicyDto,
pub token: Option<String>,
pub api_key: Option<String>,
pub user_id: Option<String>,
pub account_token: Option<String>,
pub server_id: Option<String>,
pub machine_id: Option<String>,
pub server_name: Option<String>,
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<String>,
pub enabled: bool,
pub options: Option<ConfigInputOptions>,
pub media_server: Option<MediaServerInputConfig>,
pub aliases: Option<Vec<ConfigInputAlias>>,
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<ConfigProvider>]) -> Result<Option<PathBuf>, 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 {
+1
View File
@@ -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;
@@ -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,
),
}
};
@@ -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));
@@ -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![];
+54 -4
View File
@@ -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<AppConfig>, 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<str>) -> 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<AppConfig>,
file_path: &Path,
playlist: Vec<PlaylistGroup>,
) -> (Vec<PlaylistGroup>, Result<(), TuliproxError>) {
persist_input_library_playlist(app_config, file_path, playlist).await
}
pub async fn load_input_media_server_playlist(
app_config: &Arc<AppConfig>,
file_path: &Path,
) -> Result<Vec<PlaylistGroup>, 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()
@@ -436,6 +436,7 @@ macro_rules! impl_single_file_disk_source {
impl_single_file_disk_source!(M3u, Arc<str>, M3uPlaylistItem);
impl_single_file_disk_source!(LocalLibrary, UUIDType, XtreamPlaylistItem);
impl_single_file_disk_source!(MediaServer, UUIDType, XtreamPlaylistItem);
pub struct MemoryPlaylistSource {
playlist: Arc<Vec<PlaylistGroup>>,
+3
View File
@@ -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",
+2 -3
View File
@@ -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::*;
@@ -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! {
@@ -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! {
<div class="tp__staged-input-view">
@@ -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<InputType> 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,
@@ -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);
}
}
}
@@ -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 {
@@ -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];
+1 -4
View File
@@ -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", &"<redacted>")
.finish()
f.debug_struct("SetupWebUserCredentialDto").field("username", &self.username).field("password", &"***").finish()
}
}
+705 -3
View File
@@ -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<str>) -> 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<String>,
#[serde(default, skip_serializing_if = "is_blank_optional_string")]
pub key: Option<String>,
#[serde(default, skip_serializing_if = "is_blank_optional_string")]
pub name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub kind: Option<MediaServerLibraryKindDto>,
}
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<MediaServerLibrarySelectorDto>,
#[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<String>,
#[serde(default, skip_serializing_if = "is_blank_optional_string")]
pub api_key: Option<String>,
#[serde(default, skip_serializing_if = "is_blank_optional_string")]
pub user_id: Option<String>,
#[serde(default, skip_serializing_if = "is_blank_optional_string")]
pub account_token: Option<String>,
#[serde(default, skip_serializing_if = "is_blank_optional_string")]
pub server_id: Option<String>,
#[serde(default, skip_serializing_if = "is_blank_optional_string")]
pub machine_id: Option<String>,
#[serde(default, skip_serializing_if = "is_blank_optional_string")]
pub server_name: Option<String>,
#[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<str>) -> 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<ConfigInputOptionsDto>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub media_server: Option<MediaServerInputConfigDto>,
#[serde(default, skip_serializing_if = "is_blank_optional_string")]
pub cache_duration: Option<String>,
#[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<u16, TuliproxError> {
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 {
+24 -30
View File
@@ -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 => {}
}
};
}
+2 -2
View File
@@ -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", &"<redacted>")
.field("password", &"<redacted>")
.field("username", &"***")
.field("password", &"***")
.finish()
}
}
+2
View File
@@ -22,6 +22,8 @@ pub fn is_blank_optional_string(s: &Option<String>) -> bool {
s.as_ref().is_none_or(|s| s.chars().all(|c| c.is_whitespace()))
}
pub fn is_non_blank_optional_string(s: &Option<String>) -> bool { !is_blank_optional_string(s) }
pub fn is_blank_optional_arc_str(s: &Option<Arc<str>>) -> bool {
s.as_ref().is_none_or(|s| s.chars().all(|c| c.is_whitespace()))
}