mirror of
https://github.com/euzu/tuliprox.git
synced 2026-10-03 14:32:08 +02:00
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.
This commit is contained in:
@@ -1,5 +1,6 @@
|
||||
# **tuliprox** — Self-Hosted Media Gateway & Playlist Processor
|
||||
|
||||

|
||||
`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.
|
||||
|
||||

|
||||
|
||||
## 🏆 Key Features
|
||||
|
||||
### 1. Fast, Lightweight, and Built for 24/7 Operation
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
};
|
||||
|
||||
@@ -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::<str>::from("provider"),
|
||||
stream_url: Arc::<str>::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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<Arc<str>>,
|
||||
media_started: Option<Arc<AtomicBool>>,
|
||||
/// Guards against emitting the confirmation more than once per stream.
|
||||
lease_confirmed: bool,
|
||||
lease_request_id: Option<tuliprox_core::model::PlaybackRequestId>,
|
||||
@@ -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<Arc<str>> = session_token
|
||||
.filter(|_| stream_details.provider_handle.is_some() || stream_details.has_deferred_provider_open())
|
||||
.map(Arc::<str>::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,
|
||||
|
||||
@@ -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<Arc<str>, Option<Arc<str>>> = HashMap::new();
|
||||
for item in &playlist_items {
|
||||
let Some(identity) = m3u_stream_url_identity(&item.url) else { continue };
|
||||
let mut index_url = |url: &Arc<str>| {
|
||||
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");
|
||||
|
||||
@@ -691,8 +691,26 @@ impl ActiveProviderManager {
|
||||
fn get_reserved_provider_for_owner(&self, input_name: &Arc<str>, session_owner: &str) -> Option<Arc<str>> {
|
||||
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<str>) -> 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();
|
||||
|
||||
@@ -137,6 +137,8 @@ pub struct UserSession {
|
||||
pub provider: Arc<str>,
|
||||
pub stream_url: Arc<str>,
|
||||
pub provider_session_headers: HashMap<String, String>,
|
||||
/// Shared with the response body so media confirmation survives a released VOD lease.
|
||||
pub media_started: Arc<AtomicBool>,
|
||||
/// Stable suffix appended to upstream User-Agent headers for this playback session.
|
||||
pub user_agent_stream_index: Option<u64>,
|
||||
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<Arc<AtomicBool>> {
|
||||
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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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 }
|
||||
@@ -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 }
|
||||
@@ -71,7 +71,10 @@ fn marker_from_extinf(line: &str) -> Option<u32> {
|
||||
|
||||
fn name_from_extinf(line: &str) -> Option<String> {
|
||||
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<u64, TestkitError> {
|
||||
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";
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user