diff --git a/backend/app/src/api/api_utils.rs b/backend/app/src/api/api_utils.rs index a2e3aa104..6e1378c61 100644 --- a/backend/app/src/api/api_utils.rs +++ b/backend/app/src/api/api_utils.rs @@ -5268,6 +5268,10 @@ mod tests { let (manual_update_sender, _) = mpsc::channel::(1); Arc::new(AppState { + recording_capacity: crate::api::model::recording_runtime::ProviderCapacityAdapter::new( + Arc::clone(&active_provider), + Arc::clone(&connection_manager), + ), forced_targets: Arc::new(ArcSwap::from_pointee(ProcessTargets { enabled: false, inputs: Vec::new(), diff --git a/backend/app/src/api/endpoints/custom_video_stream_api.rs b/backend/app/src/api/endpoints/custom_video_stream_api.rs index c9280e2d4..b384a9b85 100644 --- a/backend/app/src/api/endpoints/custom_video_stream_api.rs +++ b/backend/app/src/api/endpoints/custom_video_stream_api.rs @@ -710,6 +710,10 @@ mod tests { let (manual_update_sender, _) = mpsc::channel::(1); Arc::new(AppState { + recording_capacity: crate::api::model::recording_runtime::ProviderCapacityAdapter::new( + Arc::clone(&active_provider), + Arc::clone(&connection_manager), + ), forced_targets: Arc::new(ArcSwap::from_pointee(crate::model::ProcessTargets { enabled: false, inputs: Vec::new(), diff --git a/backend/app/src/api/endpoints/hls_api.rs b/backend/app/src/api/endpoints/hls_api.rs index 02d2623c9..ee7f366e7 100644 --- a/backend/app/src/api/endpoints/hls_api.rs +++ b/backend/app/src/api/endpoints/hls_api.rs @@ -12143,6 +12143,10 @@ mod tests { let (manual_update_sender, _) = mpsc::channel::(1); Arc::new(AppState { + recording_capacity: crate::api::model::recording_runtime::ProviderCapacityAdapter::new( + Arc::clone(&active_provider), + Arc::clone(&connection_manager), + ), forced_targets: Arc::new(ArcSwap::from_pointee(ProcessTargets { enabled: false, inputs: Vec::new(), diff --git a/backend/app/src/api/endpoints/recording_api.rs b/backend/app/src/api/endpoints/recording_api.rs index b563fd549..4fb550dd8 100644 --- a/backend/app/src/api/endpoints/recording_api.rs +++ b/backend/app/src/api/endpoints/recording_api.rs @@ -254,8 +254,7 @@ async fn create_http_recording_task( recording_config, &app_state.recordings, &app_state.event_manager, - &app_state.active_provider, - &app_state.connection_manager, + &app_state.recording_capacity, ) .await .is_err() diff --git a/backend/app/src/api/endpoints/v1_api_playlist.rs b/backend/app/src/api/endpoints/v1_api_playlist.rs index 975f8aa63..7957a1273 100644 --- a/backend/app/src/api/endpoints/v1_api_playlist.rs +++ b/backend/app/src/api/endpoints/v1_api_playlist.rs @@ -1408,6 +1408,10 @@ mod tests { let (manual_update_sender, _) = mpsc::channel::(1); Arc::new(AppState { + recording_capacity: crate::api::model::recording_runtime::ProviderCapacityAdapter::new( + Arc::clone(&active_provider), + Arc::clone(&connection_manager), + ), forced_targets: Arc::new(ArcSwap::from_pointee(crate::model::ProcessTargets { enabled: false, inputs: Vec::new(), diff --git a/backend/app/src/api/main_api.rs b/backend/app/src/api/main_api.rs index c2cf43ab4..f22839ad7 100644 --- a/backend/app/src/api/main_api.rs +++ b/backend/app/src/api/main_api.rs @@ -345,6 +345,10 @@ async fn create_shared_data( hls_proxy, hls_provisioning: Arc::new(HlsProvisioningState::new()), active_users, + recording_capacity: crate::api::model::recording_runtime::ProviderCapacityAdapter::new( + Arc::clone(&active_provider), + Arc::clone(&connection_manager), + ), active_provider, connection_manager, event_manager, @@ -1067,6 +1071,10 @@ mod tests { let metadata_manager = Arc::new(MetadataUpdateManager::new(tokens.metadata.clone())); let (manual_update_sender, _) = mpsc::channel::(1); Arc::new(AppState { + recording_capacity: crate::api::model::recording_runtime::ProviderCapacityAdapter::new( + Arc::clone(&active_provider), + Arc::clone(&connection_manager), + ), forced_targets: Arc::new(ArcSwap::from_pointee(ProcessTargets { enabled: false, inputs: Vec::new(), diff --git a/backend/app/src/api/model/app_state.rs b/backend/app/src/api/model/app_state.rs index 525f17d70..3ddc17c71 100644 --- a/backend/app/src/api/model/app_state.rs +++ b/backend/app/src/api/model/app_state.rs @@ -434,6 +434,9 @@ pub struct AppState { pub active_users: Arc, pub active_provider: Arc, pub connection_manager: Arc, + /// Provider capacity as the DVR sees it; the adapter that keeps provider + /// details out of the recording engine. + pub recording_capacity: Arc, pub event_manager: Arc, pub cancel_tokens: Arc>, pub playlists: Arc, @@ -504,6 +507,10 @@ pub(crate) fn create_test_app_state(config: Config) -> Arc { hls_proxy: Arc::new(HlsProxyManager::new()), hls_provisioning: Arc::new(HlsProvisioningState::new()), active_users, + recording_capacity: crate::api::model::recording_runtime::ProviderCapacityAdapter::new( + Arc::clone(&active_provider), + Arc::clone(&connection_manager), + ), active_provider, connection_manager, event_manager, diff --git a/backend/app/src/api/model/app_state_view.rs b/backend/app/src/api/model/app_state_view.rs index 8b01abb24..9a722e4e5 100644 --- a/backend/app/src/api/model/app_state_view.rs +++ b/backend/app/src/api/model/app_state_view.rs @@ -71,7 +71,7 @@ macro_rules! app_state_views { app_state_views! { /// The handles the DVR needs: the recording queue and what feeds it. recording_ctx => crate::api::model::recording::recording_ctx::RecordingCtx { - app_config, recordings, event_manager, http_client, active_provider, connection_manager, + app_config, recordings, event_manager, http_client, recording_capacity, } /// The handles the HLS proxy needs: itself, plus provider allocation and diff --git a/backend/app/src/api/model/mod.rs b/backend/app/src/api/model/mod.rs index 854e9fc4d..65ba4d4ec 100644 --- a/backend/app/src/api/model/mod.rs +++ b/backend/app/src/api/model/mod.rs @@ -94,3 +94,5 @@ pub use tuliprox_session::{ meter::*, provider_dns_manager::*, provider_lineup_manager::*, qos_aggregation_manager::*, response_headers::*, stream::*, streams::*, }; + +pub mod recording_runtime; diff --git a/backend/app/src/api/model/recording_runtime.rs b/backend/app/src/api/model/recording_runtime.rs new file mode 100644 index 000000000..30dc4729f --- /dev/null +++ b/backend/app/src/api/model/recording_runtime.rs @@ -0,0 +1,42 @@ +//! Adapts the server's connection managers to the DVR's capacity port. +//! +//! Everything provider-shaped stays on this side of the boundary. The adapter +//! does not write files, choose between ffmpeg and HTTP, or change a +//! recording's state; it answers whether there is a connection to record with +//! and hands one over. That is the whole of its job, and keeping it that small +//! is what lets the DVR be tested without a provider. + +use futures::future::BoxFuture; +use std::sync::Arc; +use tokio::sync::Notify; +use tuliprox_core::model::ProviderHandle; +use tuliprox_dvr::recording::recording_capacity::{ProviderCapacity, RecordingCapacityPort}; +use tuliprox_session::{ActiveProviderManager, ConnectionManager}; + +/// The running server's provider capacity, as the DVR sees it. +pub struct ProviderCapacityAdapter { + active_provider: Arc, + connection_manager: Arc, +} + +impl ProviderCapacityAdapter { + pub fn new(active_provider: Arc, connection_manager: Arc) -> Arc { + Arc::new(Self { active_provider, connection_manager }) + } +} + +impl RecordingCapacityPort for ProviderCapacityAdapter { + fn capacities_for_input<'a>(&'a self, input_name: &'a Arc) -> BoxFuture<'a, Vec> { + Box::pin(self.active_provider.provider_capacities_for_input(input_name)) + } + + fn acquire<'a>(&'a self, input_name: &'a Arc, priority: i8) -> BoxFuture<'a, Option> { + Box::pin(self.active_provider.acquire_connection_for_download(input_name, priority)) + } + + fn release(&self, handle: Option) -> BoxFuture<'_, ()> { + Box::pin(self.connection_manager.release_provider_handle(handle)) + } + + fn capacity_changed(&self) -> Arc { self.connection_manager.capacity_notified() } +} 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 b85cced4e..0c1074747 100644 --- a/backend/app/src/api/model/streams/active_client_stream.rs +++ b/backend/app/src/api/model/streams/active_client_stream.rs @@ -1489,6 +1489,10 @@ mod tests { let (manual_update_sender, _) = mpsc::channel::(1); Arc::new(AppState { + recording_capacity: crate::api::model::recording_runtime::ProviderCapacityAdapter::new( + Arc::clone(&active_provider), + Arc::clone(&connection_manager), + ), forced_targets: Arc::new(ArcSwap::from_pointee(ProcessTargets { enabled: false, inputs: Vec::new(), @@ -1560,6 +1564,10 @@ mod tests { let (manual_update_sender, _) = mpsc::channel::(1); Arc::new(AppState { + recording_capacity: crate::api::model::recording_runtime::ProviderCapacityAdapter::new( + Arc::clone(&active_provider), + Arc::clone(&connection_manager), + ), forced_targets: Arc::new(ArcSwap::from_pointee(ProcessTargets { enabled: false, inputs: Vec::new(), diff --git a/backend/dvr/src/recording/mod.rs b/backend/dvr/src/recording/mod.rs index 01e11ed68..a3dfe93b7 100644 --- a/backend/dvr/src/recording/mod.rs +++ b/backend/dvr/src/recording/mod.rs @@ -1,3 +1,4 @@ +pub mod recording_capacity; pub mod recording_catalog_access; pub mod recording_conflict; pub mod recording_ctx; diff --git a/backend/dvr/src/recording/recording_capacity.rs b/backend/dvr/src/recording/recording_capacity.rs new file mode 100644 index 000000000..f736ed5f1 --- /dev/null +++ b/backend/dvr/src/recording/recording_capacity.rs @@ -0,0 +1,135 @@ +//! What the DVR needs from whatever owns provider connections. +//! +//! The DVR decides whether to record, what to write and when to stop. It does +//! not decide whether there is a connection available to do it with, and it +//! must not need a real provider to be tested. This is the seam: the app +//! adapts its connection managers to it, and tests supply a stand-in. + +use futures::future::BoxFuture; +use std::sync::Arc; +use tokio::sync::Notify; +use tuliprox_core::model::ProviderHandle; + +/// One provider's capacity for an input: name, connections in use, and limit. +/// +/// A limit of zero means unlimited. +pub type ProviderCapacity = (Arc, usize, usize); + +/// Provider capacity, as the recording worker sees it. +/// +/// Object-safe on purpose -- `RecordingCtx` holds one of these behind `dyn`, so +/// the whole DVR does not become generic over its provider. +pub trait RecordingCapacityPort: Send + Sync { + /// What every provider serving this input currently looks like. + fn capacities_for_input<'a>(&'a self, input_name: &'a Arc) -> BoxFuture<'a, Vec>; + + /// Take a connection slot, or `None` when none is free. + /// + /// Lower priority values are stronger, matching the queue. + fn acquire<'a>(&'a self, input_name: &'a Arc, priority: i8) -> BoxFuture<'a, Option>; + + /// Give a slot back. Accepts `None` so callers can release unconditionally + /// on paths where they may never have held one. + fn release(&self, handle: Option) -> BoxFuture<'_, ()>; + + /// Fires when a connection is freed anywhere, so waiters can re-check. + fn capacity_changed(&self) -> Arc; +} + +#[cfg(test)] +pub mod stub { + //! A capacity port that answers from a script, for testing the worker + //! without a provider. + + use super::{ProviderCapacity, RecordingCapacityPort}; + use futures::future::BoxFuture; + use std::sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, Mutex, + }; + use tokio::sync::Notify; + use tuliprox_core::model::ProviderHandle; + + /// Counts what the worker asked for, and hands out slots on request. + pub struct StubCapacity { + capacities: Mutex>, + /// `None` means "no slot free", which is what makes a worker wait. + grants: Mutex>>, + notify: Arc, + pub acquires: AtomicUsize, + pub releases: AtomicUsize, + } + + impl StubCapacity { + fn with_capacity(in_use: usize, limit: usize) -> Arc { + Arc::new(Self { + capacities: Mutex::new(vec![(Arc::from("provider"), in_use, limit)]), + grants: Mutex::new(Vec::new()), + notify: Arc::new(Notify::new()), + acquires: AtomicUsize::new(0), + releases: AtomicUsize::new(0), + }) + } + + /// A provider with room. + pub fn with_room() -> Arc { Self::with_capacity(0, 4) } + + /// A provider that is full, so every request has to wait. + pub fn full() -> Arc { Self::with_capacity(4, 4) } + + pub fn acquire_count(&self) -> usize { self.acquires.load(Ordering::SeqCst) } + + pub fn release_count(&self) -> usize { self.releases.load(Ordering::SeqCst) } + } + + impl RecordingCapacityPort for StubCapacity { + fn capacities_for_input<'a>(&'a self, _input_name: &'a Arc) -> BoxFuture<'a, Vec> { + Box::pin(async move { self.capacities.lock().expect("capacities").clone() }) + } + + fn acquire<'a>(&'a self, _input_name: &'a Arc, _priority: i8) -> BoxFuture<'a, Option> { + Box::pin(async move { + self.acquires.fetch_add(1, Ordering::SeqCst); + self.grants.lock().expect("grants").pop().flatten() + }) + } + + fn release(&self, _handle: Option) -> BoxFuture<'_, ()> { + Box::pin(async move { + self.releases.fetch_add(1, Ordering::SeqCst); + }) + } + + fn capacity_changed(&self) -> Arc { Arc::clone(&self.notify) } + } +} + +#[cfg(test)] +mod tests { + use super::{stub::StubCapacity, RecordingCapacityPort}; + use std::sync::Arc; + + #[tokio::test] + async fn a_full_provider_offers_no_room_and_grants_nothing() { + // This is what makes a worker wait, and until now it could only be + // produced by a real provider actually being full. + let capacity = StubCapacity::full(); + let input: Arc = Arc::from("provider"); + + let reported = capacity.capacities_for_input(&input).await; + assert_eq!(reported, vec![(Arc::from("provider"), 4, 4)], "no headroom"); + assert!(capacity.acquire(&input, 0).await.is_none(), "and nothing to hand out"); + assert_eq!(capacity.acquire_count(), 1, "the attempt is observable, which is the point"); + } + + #[tokio::test] + async fn releasing_is_counted_even_when_no_slot_was_held() { + // The worker releases unconditionally on paths where it may never have + // acquired. That has to be harmless, and countable, so a leak or a + // double release is something a test can see. + let capacity = StubCapacity::with_room(); + capacity.release(None).await; + capacity.release(None).await; + assert_eq!(capacity.release_count(), 2); + } +} diff --git a/backend/dvr/src/recording/recording_ctx.rs b/backend/dvr/src/recording/recording_ctx.rs index 56f234a87..c827e06df 100644 --- a/backend/dvr/src/recording/recording_ctx.rs +++ b/backend/dvr/src/recording/recording_ctx.rs @@ -5,12 +5,12 @@ //! four directly keeps the DVR independent of the shape of the server's root //! state, which is what lets it live outside `api`. -use crate::recording::recording_queue::RecordingQueue; +use crate::recording::{recording_capacity::RecordingCapacityPort, recording_queue::RecordingQueue}; use arc_swap::ArcSwap; use reqwest::Client; use std::sync::Arc; use tuliprox_core::model::AppConfig; -use tuliprox_session::{ActiveProviderManager, ConnectionManager, EventManager}; +use tuliprox_session::EventManager; /// Everything the DVR reads from the running server. #[derive(Clone)] @@ -23,11 +23,10 @@ pub struct RecordingCtx { pub event_manager: Arc, /// Shared HTTP client, swapped when the proxy configuration changes. pub http_client: Arc>, - /// Provider capacity, used to acquire a connection slot before a - /// recording starts and to release it when the task ends. - pub active_provider: Arc, - /// Connection registry; its capacity signal wakes waiting recordings. - pub connection_manager: Arc, + /// Provider capacity: acquiring a connection slot before a recording + /// starts, releasing it when the task ends, and the signal that wakes + /// waiters when one is freed. + pub recording_capacity: Arc, } impl RecordingCtx { diff --git a/backend/dvr/src/recording/recording_transfer.rs b/backend/dvr/src/recording/recording_transfer.rs index 626a83eb7..2cbff1e57 100644 --- a/backend/dvr/src/recording/recording_transfer.rs +++ b/backend/dvr/src/recording/recording_transfer.rs @@ -12,6 +12,7 @@ use crate::{ validate_resume_response, ResponseSnapshot, ResumeValidationError, ResumeValidator, }, recording::{ + recording_capacity::RecordingCapacityPort, recording_ctx::RecordingCtx, recording_notification::LifecycleEvent, recording_notification_adapter::{build_marker, decide, message_for, DispatchDecision}, @@ -44,7 +45,7 @@ use tuliprox_core::{ utils::{async_file_writer, request, request::create_client, IO_BUFFER_SIZE}, }; use tuliprox_messaging::send_message; -use tuliprox_session::{ActiveProviderManager, ConnectionManager, EventManager, EventMessage}; +use tuliprox_session::{EventManager, EventMessage}; const DOWNLOAD_PROGRESS_LOG_INTERVAL: Duration = Duration::from_secs(5); const DOWNLOAD_PROGRESS_LOG_BYTES: u64 = 16 * 1024 * 1024; @@ -1204,8 +1205,7 @@ pub async fn ensure_recording_worker_running( download_cfg: &RecordingConfig, download_queue: &Arc, event_manager: &Arc, - active_provider: &Arc, - connection_manager: &Arc, + capacity: &Arc, ) -> Result<(), String> { let mut worker_running = download_queue.worker_running.write().await; if *worker_running { @@ -1240,8 +1240,7 @@ pub async fn ensure_recording_worker_running( let control_signal = Arc::clone(&dq.control_signal); let control_notify = Arc::clone(&dq.control_notify); let event_manager = Arc::clone(event_manager); - let active_provider = Arc::clone(active_provider); - let connection_manager = Arc::clone(connection_manager); + let capacity = Arc::clone(capacity); let download_cfg = download_cfg.clone(); let app_config = Arc::new(cfg.clone()); @@ -1271,7 +1270,7 @@ pub async fn ensure_recording_worker_running( let window_deadline = dq.active.read().await.as_ref().and_then(recording_deadline_instant); if let Some(input_name) = input_name { loop { - let capacities = active_provider.provider_capacities_for_input(&input_name).await; + let capacities = capacity.capacities_for_input(&input_name).await; if background_download_should_wait(priority, &capacities, &download_cfg) { if let Err(err) = broadcast_worker_mutation( &event_manager, @@ -1310,9 +1309,7 @@ pub async fn ensure_recording_worker_running( } continue; } - if let Some(handle) = - active_provider.acquire_connection_for_download(&input_name, priority).await - { + if let Some(handle) = capacity.acquire(&input_name, priority).await { break ProviderAcquireResult::Acquired(Some(handle)); } if *control_signal.read().await == RecordingControl::Cancel { @@ -1379,12 +1376,12 @@ pub async fn ensure_recording_worker_running( publish_recording_change(&event_manager); } Ok(None) => { - connection_manager.release_provider_handle(handle).await; + capacity.release(handle).await; error!("Download worker active task changed after provider acquire"); break 'worker; } Err(err) => { - connection_manager.release_provider_handle(handle).await; + capacity.release(handle).await; error!("Download worker commit failed after provider acquire: {err}"); break 'worker; } @@ -1440,7 +1437,7 @@ pub async fn ensure_recording_worker_running( let execution_result = { let Some(download) = active_download_snapshot_for_worker(&dq.active, &worker_uuid).await else { - connection_manager.release_provider_handle(provider_handle).await; + capacity.release(provider_handle).await; break 'worker; }; match download.kind { @@ -1525,7 +1522,7 @@ pub async fn ensure_recording_worker_running( match execution_result { DownloadExecutionResult::Completed => { - connection_manager.release_provider_handle(provider_handle).await; + capacity.release(provider_handle).await; let measured_bytes = { let active = dq.active.read().await; match active.as_ref() { @@ -1577,7 +1574,7 @@ pub async fn ensure_recording_worker_running( } } DownloadExecutionResult::Paused => { - connection_manager.release_provider_handle(provider_handle).await; + capacity.release(provider_handle).await; if let Err(err) = broadcast_required_worker_mutation( &event_manager, set_active_download_state( @@ -1596,7 +1593,7 @@ pub async fn ensure_recording_worker_running( break; } DownloadExecutionResult::Cancelled => { - connection_manager.release_provider_handle(provider_handle).await; + capacity.release(provider_handle).await; if let Err(err) = broadcast_required_worker_mutation( &event_manager, cancel_active_and_promote(&dq, &worker_uuid).await, @@ -1607,7 +1604,7 @@ pub async fn ensure_recording_worker_running( } } DownloadExecutionResult::Preempted => { - connection_manager.release_provider_handle(provider_handle).await; + capacity.release(provider_handle).await; let control = *control_signal.read().await; match control { RecordingControl::Restart => warn!( @@ -1636,7 +1633,7 @@ pub async fn ensure_recording_worker_running( } } DownloadExecutionResult::Retryable(_err) => { - connection_manager.release_provider_handle(provider_handle).await; + capacity.release(provider_handle).await; warn!("Retrying active download after transient failure"); let retry_commit = prepare_active_retry(&dq, &worker_uuid, &download_cfg).await; let retry_delay_secs = match retry_commit { @@ -1739,7 +1736,7 @@ pub async fn ensure_recording_worker_running( } } DownloadExecutionResult::Failed(err) => { - connection_manager.release_provider_handle(provider_handle).await; + capacity.release(provider_handle).await; warn!("Download failed permanently: {err}"); let committed = finish_active_and_promote(&dq, &worker_uuid, |fd| { fd.finished = true; @@ -1800,8 +1797,7 @@ pub fn spawn_recording_services(ctx: &RecordingCtx, cancel_token: &CancellationT recording_cfg, &ctx.recordings, Arc::clone(&ctx.event_manager), - Arc::clone(&ctx.active_provider), - Arc::clone(&ctx.connection_manager), + Arc::clone(&ctx.recording_capacity), cancel_token.clone(), ); } @@ -1819,8 +1815,7 @@ pub async fn resume_recording_worker_if_needed( recording_cfg, &ctx.recordings, &ctx.event_manager, - &ctx.active_provider, - &ctx.connection_manager, + &ctx.recording_capacity, ) .await } @@ -1830,13 +1825,12 @@ fn start_recording_scheduler( recording_cfg: RecordingConfig, recordings: &Arc, event_manager: Arc, - active_provider: Arc, - connection_manager: Arc, + capacity: Arc, cancel_token: CancellationToken, ) { - let capacity_notify = connection_manager.capacity_notified(); + let capacity_notify = capacity.capacity_changed(); let slot_waiters = Arc::clone(&recordings.slot_waiters); - let bridge_active_provider = Arc::clone(&active_provider); + let bridge_capacity = Arc::clone(&capacity); let bridge_recording_cfg = recording_cfg.clone(); let bridge_cancel_token = cancel_token.clone(); tokio::spawn(async move { @@ -1857,7 +1851,7 @@ fn start_recording_scheduler( let capacities = if let Some(capacities) = capacities_by_input.get(input_name) { capacities.clone() } else { - let capacities = bridge_active_provider.provider_capacities_for_input(input_name).await; + let capacities = bridge_capacity.capacities_for_input(input_name).await; capacities_by_input.insert(Arc::clone(input_name), capacities.clone()); capacities }; @@ -1895,8 +1889,7 @@ fn start_recording_scheduler( &recording_cfg, &scheduler_recordings, &event_manager, - &active_provider, - &connection_manager, + &capacity, ) .await; }