Fix M3U alias credential rewriting on provider fallback (#852)

* Fix M3U alias credential rewriting on provider fallback

  Rewrite opaque authentication query parameters when stream allocation
  switches from the primary M3U provider to an alias.

  Support unambiguous cross-key mappings such as token to api_key while
  preserving unrelated query parameters and rejecting ambiguous mappings.

  Add regression coverage for URL rewriting and provider allocation.
This commit is contained in:
euzu
2026-09-01 15:35:57 +02:00
committed by GitHub
parent a3c97c7754
commit 331a869859
7 changed files with 513 additions and 31 deletions
@@ -642,6 +642,11 @@ pub(crate) async fn download_input<E: EventSink + Clone + 'static, M: MetadataUp
}
}
if input.input_type == InputType::M3u {
let alias_errors = download_m3u_alias_playlists(ctx, input).await;
playlist_download_result.download_err.extend(alias_errors);
}
if mark_as_processed && !playlist_download_result.partial && error.is_none() && !playlist.is_empty() {
// Mark after persist/load so other workers only see this input as ready when data is usable.
ctx.mark_input_downloaded(input.name.clone()).await;
@@ -653,6 +658,45 @@ pub(crate) async fn download_input<E: EventSink + Clone + 'static, M: MetadataUp
(playlist_download_result.download_err, playlist, error, playlist_download_result.partial)
}
async fn download_m3u_alias_playlists<E: EventSink + Clone + 'static, M: MetadataUpdateSink>(
ctx: &PlaylistProcessingContext<E, M>,
input: &ConfigInput,
) -> Vec<TuliproxError> {
let Some(aliases) = input.get_enabled_aliases() else { return vec![] };
let mut errors = Vec::new();
for alias in aliases {
if ctx.is_input_downloaded(&alias.name).await {
continue;
}
let mut alias_input = input.as_input(alias);
// A user-provided raw-playlist persist path belongs to the primary input. Alias
// snapshots use their own internal storage so accounts never overwrite each other.
alias_input.persist = None;
alias_input.epg = None;
let alias_input = Arc::new(alias_input);
let (mut alias_errors, mut alias_playlist, storage_error, partial) =
Box::pin(download_input(ctx, &alias_input, false)).await;
let alias_had_errors = !alias_errors.is_empty() || storage_error.is_some();
errors.append(&mut alias_errors);
if let Some(storage_error) = storage_error {
errors.push(storage_error);
}
if partial {
errors.push(TuliproxError::RepositoryPlaylist(format!(
"M3U alias '{}' returned a partial playlist",
alias.name
)));
} else if alias_playlist.is_empty() && !alias_had_errors {
errors.push(TuliproxError::RepositoryPlaylist(format!("M3U alias '{}' playlist is empty", alias.name)));
}
}
errors
}
pub(crate) fn create_broadcast_callback<E: EventSink + Clone + 'static>(events: &E) -> StepMeasureCallback {
let events = events.clone();
Box::new(move |context: &str, msg: &str| {
@@ -9,7 +9,7 @@ use shared::{
},
utils::Internable,
};
use tuliprox_core::model::{CompiledMappingRule, CompiledTargetMappings, Config};
use tuliprox_core::model::{CompiledMappingRule, CompiledTargetMappings, Config, ConfigInputAlias};
fn serialize_without_trailing_fields<T: serde::Serialize>(value: &T, trailing_fields: &[u8]) -> Vec<u8> {
let mut encoded = rmp_serde::to_vec(value).expect("playlist item should serialize");
@@ -896,6 +896,137 @@ mod mapping_stage {
}
}
#[tokio::test]
async fn m3u_alias_playlist_is_downloaded_and_indexed_separately() {
let temp = tempfile::tempdir().expect("temp dir should be created");
let primary_playlist_path = temp.path().join("primary.m3u");
let alias_playlist_path = temp.path().join("backup.m3u");
tokio::fs::write(
&primary_playlist_path,
"#EXTM3U\n#EXTINF:-1 tvg-id=\"323\",Channel\nhttp://stream.example:4000/323/mono.m3u8?token=primary-stream-token\n",
)
.await
.expect("primary fixture should be written");
tokio::fs::write(
&alias_playlist_path,
"#EXTM3U\n#EXTINF:-1 tvg-id=\"323\",Channel\nhttp://stream.example:4000/323/mono.m3u8?token=backup-stream-token\n",
)
.await
.expect("alias fixture should be written");
let ctx = processing_context();
let config =
Config { storage_dir: temp.path().join("storage").to_string_lossy().into_owned(), ..Config::default() };
ctx.config.config.store(Arc::new(config));
let input = Arc::new(ConfigInput {
id: 1,
name: "primary-account".intern(),
input_type: InputType::M3u,
url: primary_playlist_path.to_string_lossy().into_owned(),
enabled: true,
aliases: Some(vec![ConfigInputAlias {
id: 2,
name: "backup-account".intern(),
url: alias_playlist_path.to_string_lossy().into_owned(),
username: None,
password: None,
priority: 1,
max_connections: 1,
exp_date: None,
enabled: true,
stalker: None,
}]),
..ConfigInput::default()
});
let (errors, mut primary_playlist, storage_error, partial) = download_input(&ctx, &input, false).await;
assert!(errors.is_empty(), "unexpected download errors: {errors:?}");
assert!(storage_error.is_none(), "unexpected primary storage error: {storage_error:?}");
assert!(!partial);
assert!(!primary_playlist.is_empty());
let alias_url = tuliprox_repository::load_input_m3u_stream_url(
&ctx.config,
&"backup-account".intern(),
"http://stream.example:4000/323/mono.m3u8?token=primary-stream-token",
)
.await
.expect("alias URL lookup should succeed");
assert_eq!(alias_url.as_deref(), Some("http://stream.example:4000/323/mono.m3u8?token=backup-stream-token"));
}
#[tokio::test]
async fn failed_m3u_alias_is_retried_after_primary_input_is_processed() {
let temp = tempfile::tempdir().expect("temp dir should be created");
let primary_playlist_path = temp.path().join("main.m3u");
let alias_playlist_path = temp.path().join("retry.m3u");
tokio::fs::write(
&primary_playlist_path,
"#EXTM3U\n#EXTINF:-1 tvg-id=\"323\",Channel\nhttp://stream.example:4000/323/mono.m3u8?token=main-stream-token\n",
)
.await
.expect("primary fixture should be written");
let ctx = processing_context();
let config =
Config { storage_dir: temp.path().join("storage").to_string_lossy().into_owned(), ..Config::default() };
ctx.config.config.store(Arc::new(config));
let input = Arc::new(ConfigInput {
id: 1,
name: "main-account".intern(),
input_type: InputType::M3u,
url: primary_playlist_path.to_string_lossy().into_owned(),
enabled: true,
aliases: Some(vec![ConfigInputAlias {
id: 2,
name: "retry-account".intern(),
url: alias_playlist_path.to_string_lossy().into_owned(),
username: None,
password: None,
priority: 1,
max_connections: 1,
exp_date: None,
enabled: true,
stalker: None,
}]),
..ConfigInput::default()
});
let (first_errors, mut first_playlist, first_storage_error, first_partial) =
download_input(&ctx, &input, false).await;
assert!(!first_errors.is_empty(), "missing alias should report an error");
assert!(first_storage_error.is_none(), "primary storage should succeed: {first_storage_error:?}");
assert!(!first_partial);
assert!(!first_playlist.is_empty());
assert!(ctx.is_input_downloaded("main-account").await);
assert!(!ctx.is_input_downloaded("retry-account").await);
tokio::fs::write(
&alias_playlist_path,
"#EXTM3U\n#EXTINF:-1 tvg-id=\"323\",Channel\nhttp://stream.example:4000/323/mono.m3u8?token=retry-stream-token\n",
)
.await
.expect("alias fixture should be written");
let (second_errors, mut second_playlist, second_storage_error, second_partial) =
download_input(&ctx, &input, false).await;
assert!(second_errors.is_empty(), "unexpected retry errors: {second_errors:?}");
assert!(second_storage_error.is_none(), "primary storage should remain readable: {second_storage_error:?}");
assert!(!second_partial);
assert!(!second_playlist.is_empty());
assert!(ctx.is_input_downloaded("retry-account").await);
let alias_url = tuliprox_repository::load_input_m3u_stream_url(
&ctx.config,
&"retry-account".intern(),
"http://stream.example:4000/323/mono.m3u8?token=main-stream-token",
)
.await
.expect("retried alias URL lookup should succeed");
assert_eq!(alias_url.as_deref(), Some("http://stream.example:4000/323/mono.m3u8?token=retry-stream-token"));
}
#[test]
fn persist_filter_runs_after_after_epg_mapping() {
let runtime = Runtime::new().expect("runtime");