From 305587f40bbfd29c52c5bcd78805b1dcee16e1fc Mon Sep 17 00:00:00 2001 From: euzu <33094714+euzu@users.noreply.github.com> Date: Sun, 27 Sep 2026 10:23:35 +0200 Subject: [PATCH] fix: retry unstarted streams on an available provider alias (#885) * **Bug Fixes** * Playback retries can use another available provider when playback has not started, helping requests proceed if the initially selected provider is unavailable. * Once playback has started, retries continue using the same eligible provider to preserve playback continuity. * Provider selection follows the same rules for regular playback and forced reopens, reducing unnecessary provider pinning before media begins. * Series episode URLs can resolve to the correct account-specific playback address. --- README.md | 19 +++-- backend/app/src/api/api_utils/mod.rs | 75 ++++++++++++++---- backend/app/src/api/api_utils/tests.rs | 13 ++++ .../app/src/api/endpoints/hls_api/tests.rs | 1 + .../api/model/streams/active_client_stream.rs | 14 +++- backend/repository/src/m3u_repository.rs | 72 ++++++++++++++++- .../session/src/active_provider_manager.rs | 70 ++++++++++++++++- .../session/src/active_user_manager/mod.rs | 14 ++++ .../session/src/active_user_manager/tests.rs | 5 ++ .../provider-pool-unstarted-series-retry.yml | 77 +++++++++++++++++++ .../provider-pool-unstarted-vod-retry.yml | 77 +++++++++++++++++++ tools/tuliprox-testkit/src/discovery.rs | 15 +++- tools/tuliprox-testkit/src/main.rs | 5 ++ 13 files changed, 422 insertions(+), 35 deletions(-) create mode 100644 test/fixtures/testkit/scenarios/provider-pool-unstarted-series-retry.yml create mode 100644 test/fixtures/testkit/scenarios/provider-pool-unstarted-vod-retry.yml diff --git a/README.md b/README.md index 2b3f97a2d..4a5f285d0 100644 --- a/README.md +++ b/README.md @@ -1,5 +1,6 @@ # **tuliprox** — Self-Hosted Media Gateway & Playlist Processor +![tuliprox logo](https://github.com/user-attachments/assets/8ef9ea79-62ff-4298-978f-22326c5c3d02) `tuliprox` is a high-performance, self-hosted media gateway for bringing IPTV providers, media servers, playlists, EPG data, and local media libraries together behind one clean and controllable interface. @@ -9,6 +10,14 @@ and publishes the result in the formats your clients already understand. **One service. Multiple sources. Multiple users. Multiple output formats. Full control.** +> **Legal Notice** +> +> `tuliprox` does not provide, host, sell, or distribute media content or access credentials. +> Users are solely responsible for ensuring that they have the necessary rights and authorization to access, process, +> proxy, record, or redistribute any media sources configured with `tuliprox`. +> +> The software is intended for use with legally obtained and properly authorized content and services. + ## ✨ Why Tuliprox? Tuliprox is built for more than playlist conversion. It is designed to become the central media gateway for a self-hosted setup: @@ -40,16 +49,6 @@ Tuliprox is built for more than playlist conversion. It is designed to become th | **Storage** | Embedded B+Tree engine, WAL protection, mmap scans, compression, compaction — no external DB required | | **Deployment** | Docker, Docker Compose templates, Raspberry Pi, NAS, VPS, x86 and ARM | -> **Legal Notice** -> -> `tuliprox` does not provide, host, sell, or distribute media content or access credentials. -> Users are solely responsible for ensuring that they have the necessary rights and authorization to access, process, -> proxy, record, or redistribute any media sources configured with `tuliprox`. -> -> The software is intended for use with legally obtained and properly authorized content and services. - -![tuliprox logo](https://github.com/user-attachments/assets/8ef9ea79-62ff-4298-978f-22326c5c3d02) - ## 🏆 Key Features ### 1. Fast, Lightweight, and Built for 24/7 Operation diff --git a/backend/app/src/api/api_utils/mod.rs b/backend/app/src/api/api_utils/mod.rs index 2335cb332..ba56f1b74 100644 --- a/backend/app/src/api/api_utils/mod.rs +++ b/backend/app/src/api/api_utils/mod.rs @@ -2155,12 +2155,13 @@ pub async fn force_provider_stream_response( cleanup_forced_reopen_addrs(app_state, &user_session.token, &cleanup_addrs).await; } - // In the normal case, provider-affine playback (such as VOD, series, or catchup) must remain pinned - // to its original provider account across seeks and range reconnects. - // However, if the pinned provider account is currently exhausted or unavailable, allowing fallback - // to lineup allocation acts as an emergency failover switch ("Notfallweiche") to prevent immediate - // playback disruption when another account in the provider pool has available capacity. - let preferred_provider = Some(&user_session.provider); + // A provider stays preferred after real media flows or while its allocation is active. + // A start that produced no media may choose another alias on its next request. + // An exhausted preferred account can still fall back to the lineup. + let preferred_provider = (item_type.is_live() + || user_session.media_started.load(std::sync::atomic::Ordering::Acquire) + || app_state.active_provider.should_reuse_playback_provider(&user_session.token, &user_session.provider)) + .then_some(&user_session.provider); let allow_forced_provider_fallback = true; // Never allow provider-side grace for forced seek/session reacquire. // Over-allocation here would break provider-side one-connection limits. @@ -2188,8 +2189,8 @@ pub async fn force_provider_stream_response( connection_kind, true, Some(user_session.token.as_str()), - Some(&user_session.provider_session_headers), - true, + preferred_provider.map(|_| &user_session.provider_session_headers), + preferred_provider.is_some(), grace_mode.map(|mode| matches!(mode, crate::api::model::GraceMode::Hold)), None, ) @@ -2467,6 +2468,13 @@ pub(crate) async fn stream_response( let stream_options = get_stream_options(&app_state.app_config); let session_state = app_state.active_users.get_and_update_user_session(&user.username, session_token).await; + let pinned_provider = pinned_provider.filter(|provider| { + item_type.is_live() + || session_state + .as_ref() + .is_some_and(|session| session.media_started.load(std::sync::atomic::Ordering::Acquire)) + || app_state.active_provider.should_reuse_playback_provider(session_token, provider) + }); let mut stream_details = match create_stream_response_details( app_state, &stream_options, @@ -2663,6 +2671,35 @@ pub(crate) async fn stream_response( stream_details.shared_subscriber_id = pending_shared_cleanup.as_ref().map(tuliprox_session::PendingSharedSubscriberCleanup::capability); } + // In the no-limits path there may be no placeholder yet. The body needs the + // session's media flag before its first byte; create that session now. + let created_media_session = if !is_stream_shared + && !item_type.is_live() + && item_type.requires_provider_affinity() + && app_state.active_users.media_started_flag(&user.username, session_token).await.is_none() + { + if let Some(provider) = provider_name.as_deref() { + app_state + .active_users + .ensure_user_session_placeholder(crate::api::model::CreateUserSessionParams { + user, + session_token, + virtual_id, + provider, + stream_url: actual_request_url.as_ref(), + addr: &fingerprint.addr, + connection_permission, + connection_kind: Some(connection_kind), + socket_bound, + }) + .await; + true + } else { + false + } + } else { + false + }; let stream = match create_active_client_stream(crate::api::model::ActiveClientStreamParams { stream_details, app_state, @@ -2681,15 +2718,19 @@ pub(crate) async fn stream_response( { Ok(stream) => stream, Err(error) => { - app_state - .active_users - .release_unbound_session_reservation( - &user.username, - session_token, - activation.placeholder_transition_version, - activation.placeholder_transition_version.is_some(), - ) - .await; + if created_media_session { + app_state.active_users.terminate_session(&user.username, session_token).await; + } else { + app_state + .active_users + .release_unbound_session_reservation( + &user.username, + session_token, + activation.placeholder_transition_version, + activation.placeholder_transition_version.is_some(), + ) + .await; + } return stream_admission_rejected_response(error, &user.username); } }; diff --git a/backend/app/src/api/api_utils/tests.rs b/backend/app/src/api/api_utils/tests.rs index 325e0b0e4..0298c9fb9 100644 --- a/backend/app/src/api/api_utils/tests.rs +++ b/backend/app/src/api/api_utils/tests.rs @@ -1095,6 +1095,7 @@ async fn forced_legacy_hls_test_response( provider: Arc::clone(&input.name), stream_url: origin_url.as_str().intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(true)), user_agent_stream_index: None, addr: client_addr, socket_bound: false, @@ -1243,6 +1244,7 @@ async fn forced_reopen_stays_on_pinned_provider_account() { provider: Arc::clone(&input.name), stream_url: format!("http://{origin_a}/live/1.ts").intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(true)), user_agent_stream_index: None, addr: client_addr, socket_bound: false, @@ -1322,6 +1324,7 @@ async fn overlapping_vod_range_requests_return_correct_account_bytes() { provider: Arc::clone(&input.name), stream_url: format!("http://{origin_addr}/movie/1.ts").intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(true)), user_agent_stream_index: None, addr: client_addr, socket_bound: false, @@ -1440,6 +1443,7 @@ async fn parallel_series_range_requests_keep_both_claims_active() { provider: Arc::clone(&input.name), stream_url: format!("http://{origin_addr}/series/1.ts").intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(true)), user_agent_stream_index: None, addr: client_addr, socket_bound: false, @@ -1790,6 +1794,7 @@ fn load_test_session(input: &ConfigInput, origin_addr: SocketAddr, token: &str, provider: Arc::clone(&input.name), stream_url: format!("http://{origin_addr}/live/42.ts").intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(false)), user_agent_stream_index: None, addr, socket_bound: false, @@ -2075,6 +2080,7 @@ async fn direct_ts_eof_before_first_byte_releases_slot_without_idle_lease() { provider: Arc::clone(&input.name), stream_url: format!("http://{origin_addr}/live/42.ts").intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(false)), user_agent_stream_index: None, addr: client_addr, socket_bound: false, @@ -2164,6 +2170,7 @@ async fn direct_ts_abort_while_waiting_for_first_byte_releases_exact_request() { provider: Arc::clone(&input.name), stream_url: format!("http://{origin_addr}/live/42.ts").intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(false)), user_agent_stream_index: None, addr: client_addr, socket_bound: false, @@ -2529,6 +2536,7 @@ async fn parallel_vod_abort_preserves_sibling_claim_and_provider_stickiness() { provider: Arc::clone(&input.name), stream_url: format!("http://{origin_addr}/movie/1.ts").intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(true)), user_agent_stream_index: None, addr: client_addr, socket_bound: false, @@ -2688,6 +2696,7 @@ async fn parallel_series_abort_then_seek_reuses_account_without_stale_claim() { provider: Arc::clone(&input.name), stream_url: format!("http://{origin_addr}/series/1.ts").intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(true)), user_agent_stream_index: None, addr: client_addr, socket_bound: false, @@ -3305,6 +3314,7 @@ async fn catchup_abort_seek_and_window_change_preserve_correct_affinity() { provider: Arc::clone(&input.name), stream_url: format!("http://{origin_a}/timeshift/1.ts?window={window}").intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(true)), user_agent_stream_index: None, addr: client_addr, socket_bound: false, @@ -4397,6 +4407,7 @@ fn test_should_allow_exhausted_shared_reconnect_only_for_matching_shared_session provider: Arc::::from("provider"), stream_url: Arc::::from("http://provider/live/449924.ts"), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(false)), user_agent_stream_index: None, addr: "127.0.0.1:1234".parse().unwrap_or_else(|_| unreachable!()), socket_bound: false, @@ -5180,6 +5191,7 @@ fn create_test_session( _ => "http://provider-1.example/live/42.ts", }), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(false)), user_agent_stream_index: None, addr: "127.0.0.1:55555".parse().unwrap_or_else(|_| unreachable!()), socket_bound: item_type.uses_socket_bound_session(), @@ -9130,6 +9142,7 @@ fn session_reacquire_cleanup_addrs_excludes_current_and_deduplicates() { provider: "provider-a".intern(), stream_url: "http://localhost/movie.mkv".intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(false)), user_agent_stream_index: None, addr: seek, socket_bound: false, diff --git a/backend/app/src/api/endpoints/hls_api/tests.rs b/backend/app/src/api/endpoints/hls_api/tests.rs index 26c3933df..daa69db7d 100644 --- a/backend/app/src/api/endpoints/hls_api/tests.rs +++ b/backend/app/src/api/endpoints/hls_api/tests.rs @@ -10846,6 +10846,7 @@ fn stats_provider_test_user_session(provider: &str) -> UserSession { provider: Arc::from(provider), stream_url: Arc::from("http://origin.example.com/live/12345.m3u8"), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(false)), user_agent_stream_index: None, addr: test_addr(), socket_bound: false, diff --git a/backend/app/src/api/model/streams/active_client_stream.rs b/backend/app/src/api/model/streams/active_client_stream.rs index b14a5746e..1cb345fcd 100644 --- a/backend/app/src/api/model/streams/active_client_stream.rs +++ b/backend/app/src/api/model/streams/active_client_stream.rs @@ -27,7 +27,7 @@ use std::{ net::SocketAddr, pin::Pin, sync::{ - atomic::{AtomicU8, Ordering}, + atomic::{AtomicBool, AtomicU8, Ordering}, Arc, }, task::{Context, Poll}, @@ -265,6 +265,7 @@ struct ActiveClientStreamState { /// Playback lease owner (session token) whose provider slot is confirmed once real /// media bytes reach the client. `None` when this stream has no provider lease. lease_owner: Option>, + media_started: Option>, /// Guards against emitting the confirmation more than once per stream. lease_confirmed: bool, lease_request_id: Option, @@ -361,6 +362,9 @@ impl ActiveClientStreamState { return; }; self.lease_confirmed = true; + if let Some(media_started) = &self.media_started { + media_started.store(true, Ordering::Release); + } let request_id = self.lease_request_id.or_else(|| { self.provider_handle.as_ref().and_then(|managed| managed.handle()).and_then(|h| h.playback_request_id) }); @@ -1106,6 +1110,11 @@ pub(crate) async fn create_active_client_stream( let lease_owner: Option> = session_token .filter(|_| stream_details.provider_handle.is_some() || stream_details.has_deferred_provider_open()) .map(Arc::::from); + let media_started = if let Some(owner) = lease_owner.as_deref() { + app_state.active_users.media_started_flag(&user.username, owner).await + } else { + None + }; let owned_grace_ctx = stream_details.grace_resolution_context.clone(); let lease_request_id = stream_details .provider_handle @@ -1230,6 +1239,7 @@ pub(crate) async fn create_active_client_stream( provider_http_status: None, provider_reconnect_count: AtomicU8::new(0), lease_owner, + media_started, lease_confirmed: false, lease_request_id, request_cleanup: registered_request.into_body_cleanup(), @@ -2181,6 +2191,7 @@ mod tests { provider_http_status: None, provider_reconnect_count: AtomicU8::new(0), lease_owner: None, + media_started: None, lease_confirmed: false, lease_request_id: None, request_cleanup: None, @@ -2725,6 +2736,7 @@ mod tests { provider_http_status: None, provider_reconnect_count: AtomicU8::new(0), lease_owner: None, + media_started: None, lease_confirmed: false, lease_request_id: None, request_cleanup: None, diff --git a/backend/repository/src/m3u_repository.rs b/backend/repository/src/m3u_repository.rs index 13841fc78..9f1916365 100644 --- a/backend/repository/src/m3u_repository.rs +++ b/backend/repository/src/m3u_repository.rs @@ -414,14 +414,24 @@ pub async fn persist_input_m3u_playlist( tree.store(&m3u_path_clone).map_err(|err| cant_write_result!(RepositoryM3u, "m3u", &m3u_path_clone, err))?; let mut indexed_urls: HashMap, Option>> = HashMap::new(); - for item in &playlist_items { - let Some(identity) = m3u_stream_url_identity(&item.url) else { continue }; + let mut index_url = |url: &Arc| { + let Some(identity) = m3u_stream_url_identity(url) else { return }; if let Some(indexed_url) = indexed_urls.get_mut(identity.as_str()) { - if indexed_url.as_ref().is_some_and(|url| url.as_ref() != item.url.as_ref()) { + if indexed_url.as_ref().is_some_and(|indexed| indexed.as_ref() != url.as_ref()) { *indexed_url = None; } } else { - indexed_urls.insert(identity.into(), Some(Arc::clone(&item.url))); + indexed_urls.insert(identity.into(), Some(Arc::clone(url))); + } + }; + for item in &playlist_items { + index_url(&item.url); + if let Some(StreamProperties::Series(series)) = &item.additional_properties { + if let Some(episodes) = series.details.as_ref().and_then(|details| details.episodes.as_ref()) { + for episode in episodes { + index_url(&episode.direct_source); + } + } } } @@ -568,6 +578,7 @@ mod tests { use shared::{ model::{ ConfigPaths, LiveStreamProperties, PlaylistGroup, PlaylistItem, PlaylistItemHeader, PlaylistItemType, + SeriesStreamDetailEpisodeProperties, SeriesStreamDetailProperties, SeriesStreamProperties, StreamProperties, XtreamCluster, }, utils::Internable, @@ -685,6 +696,59 @@ mod tests { assert_eq!(properties.bitrate, 2_500_000); } + #[tokio::test] + async fn persisted_stream_url_index_resolves_series_episode_alias_token() { + let temp = tempfile::tempdir().expect("temp dir should be created"); + let app_config = test_app_config(); + let config = Config { storage_dir: temp.path().to_string_lossy().into_owned(), ..Config::default() }; + app_config.config.store(Arc::new(config)); + + let alias_name = "backup-account".intern(); + let storage_path = get_input_storage_path(&alias_name, &app_config.config.load().storage_dir) + .await + .expect("alias storage should be created"); + let playlist_path = get_input_m3u_playlist_file_path(&storage_path, &alias_name); + let playlist = vec![PlaylistGroup { + id: 1, + title: "Series".intern(), + channels: vec![PlaylistItem { + header: PlaylistItemHeader { + id: "series-1".intern(), + url: "".intern(), + item_type: PlaylistItemType::SeriesInfo, + xtream_cluster: XtreamCluster::Series, + additional_properties: Some(StreamProperties::Series(Box::new(SeriesStreamProperties { + details: Some(SeriesStreamDetailProperties::new( + None, + Vec::new(), + Some(vec![SeriesStreamDetailEpisodeProperties { + direct_source: "http://stream.example/vod/episode.mkv?token=backup-stream-token" + .intern(), + ..SeriesStreamDetailEpisodeProperties::default() + }]), + )), + ..SeriesStreamProperties::default() + }))), + ..PlaylistItemHeader::default() + }, + }], + xtream_cluster: XtreamCluster::Series, + }]; + persist_input_m3u_playlist(&app_config, &playlist_path, &playlist) + .await + .expect("alias playlist should persist"); + + let resolved = load_input_m3u_stream_url( + &app_config, + &alias_name, + "http://stream.example/vod/episode.mkv?token=primary-stream-token", + ) + .await + .expect("alias episode URL lookup should succeed"); + + assert_eq!(resolved.as_deref(), Some("http://stream.example/vod/episode.mkv?token=backup-stream-token")); + } + #[tokio::test] async fn persisted_stream_url_index_resolves_account_specific_alias_token() { let temp = tempfile::tempdir().expect("temp dir should be created"); diff --git a/backend/session/src/active_provider_manager.rs b/backend/session/src/active_provider_manager.rs index 2de48279e..b0ff22fd1 100644 --- a/backend/session/src/active_provider_manager.rs +++ b/backend/session/src/active_provider_manager.rs @@ -691,8 +691,26 @@ impl ActiveProviderManager { fn get_reserved_provider_for_owner(&self, input_name: &Arc, session_owner: &str) -> Option> { let mut leases = self.write_leases(); Self::prune_expired_leases(&mut leases); - let provider_name = leases.provider_for_owner(session_owner)?; - self.providers.is_provider_for_input(&provider_name, input_name).then_some(provider_name) + let lease = leases.lease_of_owner(session_owner)?; + let provider_name = Arc::clone(&lease.provider_name); + let confirmed = lease.state.is_confirmed(); + drop(leases); + (self.providers.is_provider_for_input(&provider_name, input_name) + && (confirmed || self.has_active_owner_for_provider(&provider_name, session_owner))) + .then_some(provider_name) + } + + /// A provisional lease only pins its provider while its allocation is active. + /// Once an unstarted request releases its slot, the next attempt may use another alias. + pub fn should_reuse_playback_provider(&self, session_owner: &str, provider_name: &Arc) -> bool { + let _transition = self.lock_capacity_transition(); + let mut leases = self.write_leases(); + Self::prune_expired_leases(&mut leases); + let confirmed = leases + .lease_of_owner(session_owner) + .is_some_and(|lease| lease.provider_name == *provider_name && lease.state.is_confirmed()); + drop(leases); + confirmed || self.has_active_owner_for_provider(provider_name, session_owner) } /// True when `session_owner` already backs a live allocation on `provider_name`, @@ -3318,6 +3336,54 @@ mod tests { manager.release_connection(&client_2_addr); } + #[tokio::test] + async fn unstarted_series_retry_uses_free_alias_while_started_playback_stays_pinned() { + let app_cfg = create_test_app_config_with_dual_provider_pool(); + let event_manager = Arc::new(EventManager::new()); + let manager = ActiveProviderManager::new(&app_cfg, &event_manager); + let input = "provider_1".intern(); + let owner = "series-playback"; + let first_addr: SocketAddr = "127.0.0.1:41010".parse().unwrap(); + let busy_addr: SocketAddr = "127.0.0.1:41011".parse().unwrap(); + let retry_addr: SocketAddr = "127.0.0.1:41012".parse().unwrap(); + + let failed_start = manager + .acquire_connection_with_lease_for_session( + &input, + &first_addr, + false, + 0, + ConnectionKind::Normal, + Some(PlaybackLeaseRef::new(owner, PlaybackKind::Series)), + ) + .expect("first series allocation"); + assert!(manager.should_reuse_playback_provider(owner, &input)); + manager.release_handle(&failed_start); + assert!(!manager.should_reuse_playback_provider(owner, &input)); + + let busy = manager + .acquire_connection(&input, &busy_addr, 0, ConnectionKind::Normal) + .expect("other playback occupies first provider"); + let retry = manager + .acquire_connection_with_lease_for_session( + &input, + &retry_addr, + false, + 0, + ConnectionKind::Normal, + Some(PlaybackLeaseRef::new(owner, PlaybackKind::Series)), + ) + .expect("retry should use free alias"); + let alias = retry.allocation.get_provider_name().expect("alias name"); + assert_eq!(alias.as_ref(), "provider_2"); + assert!(manager.should_reuse_playback_provider(owner, &alias)); + manager.refresh_adaptive_playback_lease(&alias, owner, PlaybackKind::Series, 15); + manager.confirm_identified_playback_activity(owner, retry.playback_request_id.expect("request id")); + manager.release_handle(&retry); + assert!(manager.should_reuse_playback_provider(owner, &alias)); + manager.release_handle(&busy); + } + #[tokio::test] async fn test_force_session_fallback_uses_different_provider_when_current_is_busy() { let app_cfg = create_test_app_config_with_dual_provider_pool(); diff --git a/backend/session/src/active_user_manager/mod.rs b/backend/session/src/active_user_manager/mod.rs index b8b8b2d10..10b5eef95 100644 --- a/backend/session/src/active_user_manager/mod.rs +++ b/backend/session/src/active_user_manager/mod.rs @@ -137,6 +137,8 @@ pub struct UserSession { pub provider: Arc, pub stream_url: Arc, pub provider_session_headers: HashMap, + /// Shared with the response body so media confirmation survives a released VOD lease. + pub media_started: Arc, /// Stable suffix appended to upstream User-Agent headers for this playback session. pub user_agent_stream_index: Option, pub addr: SocketAddr, @@ -1961,6 +1963,7 @@ impl ActiveUserManager { provider: params.provider.intern(), stream_url: params.stream_url.intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(AtomicBool::new(false)), user_agent_stream_index: None, addr: *params.addr, socket_bound: params.socket_bound, @@ -3094,6 +3097,17 @@ impl ActiveUserManager { self.update_user_session(username, token).await } + pub async fn media_started_flag(&self, username: &str, token: &str) -> Option> { + let users = self.connections.read().await; + users + .by_key + .get(username)? + .sessions + .iter() + .find(|session| session.token == token) + .map(|session| Arc::clone(&session.media_started)) + } + /// Session for target-scoped `virtual_id` and request token (used to recover leaked relative DVR segment paths). pub async fn find_latest_session_for_target_stream( &self, diff --git a/backend/session/src/active_user_manager/tests.rs b/backend/session/src/active_user_manager/tests.rs index 60185b7c1..df930139d 100644 --- a/backend/session/src/active_user_manager/tests.rs +++ b/backend/session/src/active_user_manager/tests.rs @@ -210,6 +210,7 @@ async fn create_user_session_normalizes_expired_lifecycle() { provider: "provider-a".intern(), stream_url: "http://localhost/live.m3u8".intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(false)), user_agent_stream_index: None, addr, socket_bound: false, @@ -272,6 +273,7 @@ async fn create_user_session_does_not_normalize_pending_provider_lifecycle() { provider: "provider-a".intern(), stream_url: "http://localhost/live.m3u8".intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(false)), user_agent_stream_index: None, addr, socket_bound: false, @@ -5885,6 +5887,7 @@ async fn check_divergence_detects_connection_count_mismatch() { provider: "provider-a".intern(), stream_url: "http://localhost/stream.ts".intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(false)), user_agent_stream_index: None, addr, socket_bound: false, @@ -5929,6 +5932,7 @@ async fn check_divergence_detects_stream_without_counted_session() { provider: "provider-a".intern(), stream_url: "http://localhost/stream.ts".intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(false)), user_agent_stream_index: None, addr, socket_bound: false, @@ -6014,6 +6018,7 @@ async fn divergence_log_rate_limited_within_cooldown_window() { provider: "provider-a".intern(), stream_url: "http://localhost/stream.ts".intern(), provider_session_headers: HashMap::new(), + media_started: Arc::new(std::sync::atomic::AtomicBool::new(false)), user_agent_stream_index: None, addr, socket_bound: false, diff --git a/test/fixtures/testkit/scenarios/provider-pool-unstarted-series-retry.yml b/test/fixtures/testkit/scenarios/provider-pool-unstarted-series-retry.yml new file mode 100644 index 000000000..7e2fa03ce --- /dev/null +++ b/test/fixtures/testkit/scenarios/provider-pool-unstarted-series-retry.yml @@ -0,0 +1,77 @@ +schema_version: 1 +name: provider-pool-unstarted-series-retry +tuliprox: + base_url: http://fixture.invalid + execution_mode: isolated_fixture + bootstrap: + command: "${env:TULIPROX_TESTKIT_SUT_BINARY}" + arguments: ["-s", "-c", "{config_file}", "-i", "{source_file}", "-a", "{api_proxy_file}"] + readiness_url: "{base_url}/healthcheck" +origin: + run_id: provider-pool-unstarted-series-retry + account_limit: 2 + limit_mode: reject_new +actors: + - id: series-viewer + agent: local + username: viewer + user_agent: tuliprox-testkit/series-viewer + client_ip: { mode: forwarded, value: 10.31.0.10 } + - id: live-blocker + agent: local + username: blocker + user_agent: tuliprox-testkit/live-blocker + client_ip: { mode: forwarded, value: 10.31.0.11 } +channels: + blocker-channel: { origin_marker: 17, protocol: xtream_ts } +steps: + - command_id: occupy-first-account + start: + actor: live-blocker + playback_id: blocker-a + session_group: blocker + channel: blocker-channel + await: { valid_frames: 3 } + assert_origin: + active_body: 1 + latest_request_account: account-a + - command_id: probe-series-on-second-account-without-media + start: + actor: series-viewer + playback_id: series-head-b + session_group: episode + vod_object: episode.mkv + method: HEAD + expected_status: 200 + assert_origin: { active_body: 1, latest_request_account: account-b } + - command_id: release-series-probe + stop: { playback_id: series-head-b } + assert_origin: { active_body: 1 } + - command_id: release-first-account + stop: { playback_id: blocker-a } + assert_origin: { active_body: 0 } + - command_id: retry-series-on-free-first-account + start: + actor: series-viewer + playback_id: series-retry-a + session_group: episode + vod_object: episode.mkv + range: "bytes=0-524287" + read_limit_bytes: 32768 + post_read_action: pause + assert_origin: + active_body: 1 + latest_request_account: account-a + no_limit_rejections: true + - command_id: release-series-retry + stop: { playback_id: series-retry-a } + assert_origin: { active_body: 0 } +policy_contract: + user_access_control: true + users: + viewer: { max_connections: 1 } + blocker: { max_connections: 1 } + admission_strategies: [] + provider_pool: + - { name: account-a, max_connections: 1 } + - { name: account-b, max_connections: 1 } diff --git a/test/fixtures/testkit/scenarios/provider-pool-unstarted-vod-retry.yml b/test/fixtures/testkit/scenarios/provider-pool-unstarted-vod-retry.yml new file mode 100644 index 000000000..0ccb507e3 --- /dev/null +++ b/test/fixtures/testkit/scenarios/provider-pool-unstarted-vod-retry.yml @@ -0,0 +1,77 @@ +schema_version: 1 +name: provider-pool-unstarted-vod-retry +tuliprox: + base_url: http://fixture.invalid + execution_mode: isolated_fixture + bootstrap: + command: "${env:TULIPROX_TESTKIT_SUT_BINARY}" + arguments: ["-s", "-c", "{config_file}", "-i", "{source_file}", "-a", "{api_proxy_file}"] + readiness_url: "{base_url}/healthcheck" +origin: + run_id: provider-pool-unstarted-vod-retry + account_limit: 2 + limit_mode: reject_new +actors: + - id: vod-viewer + agent: local + username: viewer + user_agent: tuliprox-testkit/vod-viewer + client_ip: { mode: forwarded, value: 10.31.0.10 } + - id: live-blocker + agent: local + username: blocker + user_agent: tuliprox-testkit/live-blocker + client_ip: { mode: forwarded, value: 10.31.0.11 } +channels: + blocker-channel: { origin_marker: 17, protocol: xtream_ts } +steps: + - command_id: occupy-first-account + start: + actor: live-blocker + playback_id: blocker-a + session_group: blocker + channel: blocker-channel + await: { valid_frames: 3 } + assert_origin: + active_body: 1 + latest_request_account: account-a + - command_id: probe-vod-on-second-account-without-media + start: + actor: vod-viewer + playback_id: vod-head-b + session_group: movie + vod_object: movie.mkv + method: HEAD + expected_status: 200 + assert_origin: { active_body: 1, latest_request_account: account-b } + - command_id: release-vod-probe + stop: { playback_id: vod-head-b } + assert_origin: { active_body: 1 } + - command_id: release-first-account + stop: { playback_id: blocker-a } + assert_origin: { active_body: 0 } + - command_id: retry-vod-on-free-first-account + start: + actor: vod-viewer + playback_id: vod-retry-a + session_group: movie + vod_object: movie.mkv + range: "bytes=0-524287" + read_limit_bytes: 32768 + post_read_action: pause + assert_origin: + active_body: 1 + latest_request_account: account-a + no_limit_rejections: true + - command_id: release-vod-retry + stop: { playback_id: vod-retry-a } + assert_origin: { active_body: 0 } +policy_contract: + user_access_control: true + users: + viewer: { max_connections: 1 } + blocker: { max_connections: 1 } + admission_strategies: [] + provider_pool: + - { name: account-a, max_connections: 1 } + - { name: account-b, max_connections: 1 } diff --git a/tools/tuliprox-testkit/src/discovery.rs b/tools/tuliprox-testkit/src/discovery.rs index d162eca45..5dddd2904 100644 --- a/tools/tuliprox-testkit/src/discovery.rs +++ b/tools/tuliprox-testkit/src/discovery.rs @@ -71,7 +71,10 @@ fn marker_from_extinf(line: &str) -> Option { fn name_from_extinf(line: &str) -> Option { let value = line.split("tvg-id=\"").nth(1)?.split('"').next()?; - value.strip_prefix("test-vod-").map(str::to_owned) + value.strip_prefix("test-vod-").map(str::to_owned).or_else(|| { + let title = line.split("tvg-name=\"").nth(1)?.split('"').next()?; + (title == "Test Series S01E01").then(|| "episode.mkv".to_owned()) + }) } /// Extract the numeric virtual ID from the last path segment of a URL. @@ -100,6 +103,16 @@ fn extract_virtual_id_from_url(url: &str) -> Result { mod tests { use super::*; + #[test] + fn discovers_synthetic_series_episode_without_tvg_id() { + let playlist = "#EXTM3U\n#EXTINF:-1 tvg-id=\"\" tvg-name=\"Test Series S01E01\" group-title=\"Series\",Test Series S01E01\nhttp://tuliprox/series/abc/def/7.mkv\n"; + let map = VirtualIdMap::from_m3u(playlist).expect("series episode discovery"); + assert_eq!( + map.named_playback_url("episode.mkv").expect("series episode URL"), + "http://tuliprox/series/abc/def/7.mkv" + ); + } + #[test] fn discovers_virtual_playback_urls_without_assuming_their_ids() { let playlist = "#EXTM3U\n#EXTINF:-1 tvg-id=\"test-17\",News\nhttp://tuliprox/m3u/abc/17\n#EXTINF:-1 tvg-id=\"test-vod-movie.mkv\" type=\"movie\",Movie\nhttp://tuliprox/movie/abc/def/1.mkv\n"; diff --git a/tools/tuliprox-testkit/src/main.rs b/tools/tuliprox-testkit/src/main.rs index 9efda1483..2d7940b3e 100644 --- a/tools/tuliprox-testkit/src/main.rs +++ b/tools/tuliprox-testkit/src/main.rs @@ -720,6 +720,11 @@ fn catalog_m3u(state: &OriginState, host: &str, account: Option<&str>) -> String } let _ = writeln!(catalog, "#EXTINF:-1 tvg-id=\"test-vod-movie.mkv\" tvg-type=\"movie\",Test Movie"); let _ = writeln!(catalog, "http://{host}/vod/movie.mkv?run={}{}", state.run_id.0, account_query); + let _ = writeln!( + catalog, + "#EXTINF:0 tvg-id=\"test-vod-episode.mkv\" tvg-type=\"series\" group-title=\"Test Series\",Test Series S01E01" + ); + let _ = writeln!(catalog, "http://{host}/vod/episode.mkv?run={}{}", state.run_id.0, account_query); catalog }