From bd2e0750ce0fa0d3ed8bfe1458d1ac90a71d4cc2 Mon Sep 17 00:00:00 2001 From: DarkBreakpoint Date: Wed, 2 Sep 2026 11:45:55 -0500 Subject: [PATCH] refactor(dvr): provider capacity becomes a port the app adapts to Task 13, steps 1 and 3. The recording worker held `Arc` and `Arc` directly, so the execution path could not be exercised without a real provider. Three plan items are blocked on exactly that: ffmpeg started once at the padded start and stopped once at the padded end (task 12), worker cancellation happening once (task 11), and exactly one worker start (task 16). None of them is a test-writing problem -- there was nothing to stand in for. `RecordingCapacityPort` is the seam: capacities for an input, acquire, release, and the signal that wakes waiters. Four methods, which is all the worker ever used of two managers. It is object-safe on purpose -- boxed futures rather than async fn in trait -- so `RecordingCtx` holds one behind `dyn` and the DVR does not become generic over its provider. `ProviderCapacityAdapter` in `backend/app` implements it over the real managers, which is where provider detail belongs. It does nothing else: no files, no strategy choice, no ffmpeg, no state changes. That is what step 1 asks the adapter not to do, and keeping it to four delegating methods is what makes it obvious it does not. `AppState` gains `recording_capacity` rather than the view renaming a handle, because the view macro's rule is that field names match on both sides -- a view renaming a handle would be a view lying about what it reads. `StubCapacity` is the payoff: a provider that is full, or has room, that counts what was asked of it. Its two tests are thin, but the point is that a worker can now be run against a provider that never grants a slot without waiting on a real one to be busy. Not done, and still needing Design section 9.1: steps 2 and 5, the per-field hot reload matrix. Step 4 landed earlier and step 6 already passed. --- backend/app/src/api/api_utils.rs | 4 + .../api/endpoints/custom_video_stream_api.rs | 4 + backend/app/src/api/endpoints/hls_api.rs | 4 + .../app/src/api/endpoints/recording_api.rs | 3 +- .../app/src/api/endpoints/v1_api_playlist.rs | 4 + backend/app/src/api/main_api.rs | 8 ++ backend/app/src/api/model/app_state.rs | 7 + backend/app/src/api/model/app_state_view.rs | 2 +- backend/app/src/api/model/mod.rs | 2 + .../app/src/api/model/recording_runtime.rs | 42 ++++++ .../api/model/streams/active_client_stream.rs | 8 ++ backend/dvr/src/recording/mod.rs | 1 + .../dvr/src/recording/recording_capacity.rs | 135 ++++++++++++++++++ backend/dvr/src/recording/recording_ctx.rs | 13 +- .../dvr/src/recording/recording_transfer.rs | 51 +++---- 15 files changed, 249 insertions(+), 39 deletions(-) create mode 100644 backend/app/src/api/model/recording_runtime.rs create mode 100644 backend/dvr/src/recording/recording_capacity.rs 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; }