refactor(dvr): provider capacity becomes a port the app adapts to

Task 13, steps 1 and 3. The recording worker held
`Arc<ActiveProviderManager>` and `Arc<ConnectionManager>` 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.
This commit is contained in:
DarkBreakpoint
2026-09-02 11:45:55 -05:00
parent 8776796f1e
commit bd2e0750ce
15 changed files with 249 additions and 39 deletions
+4
View File
@@ -5268,6 +5268,10 @@ mod tests {
let (manual_update_sender, _) = mpsc::channel::<crate::api::model::ManualPlaylistUpdateRequest>(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(),
@@ -710,6 +710,10 @@ mod tests {
let (manual_update_sender, _) = mpsc::channel::<crate::api::model::ManualPlaylistUpdateRequest>(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(),
+4
View File
@@ -12143,6 +12143,10 @@ mod tests {
let (manual_update_sender, _) = mpsc::channel::<ManualPlaylistUpdateRequest>(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(),
@@ -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()
@@ -1408,6 +1408,10 @@ mod tests {
let (manual_update_sender, _) = mpsc::channel::<crate::api::model::ManualPlaylistUpdateRequest>(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(),
+8
View File
@@ -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::<ManualPlaylistUpdateRequest>(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(),
+7
View File
@@ -434,6 +434,9 @@ pub struct AppState {
pub active_users: Arc<ActiveUserManager>,
pub active_provider: Arc<ActiveProviderManager>,
pub connection_manager: Arc<ConnectionManager>,
/// Provider capacity as the DVR sees it; the adapter that keeps provider
/// details out of the recording engine.
pub recording_capacity: Arc<dyn tuliprox_dvr::recording::recording_capacity::RecordingCapacityPort>,
pub event_manager: Arc<EventManager>,
pub cancel_tokens: Arc<ArcSwap<CancelTokens>>,
pub playlists: Arc<PlaylistStorageState>,
@@ -504,6 +507,10 @@ pub(crate) fn create_test_app_state(config: Config) -> Arc<AppState> {
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,
+1 -1
View File
@@ -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
+2
View File
@@ -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;
@@ -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<ActiveProviderManager>,
connection_manager: Arc<ConnectionManager>,
}
impl ProviderCapacityAdapter {
pub fn new(active_provider: Arc<ActiveProviderManager>, connection_manager: Arc<ConnectionManager>) -> Arc<Self> {
Arc::new(Self { active_provider, connection_manager })
}
}
impl RecordingCapacityPort for ProviderCapacityAdapter {
fn capacities_for_input<'a>(&'a self, input_name: &'a Arc<str>) -> BoxFuture<'a, Vec<ProviderCapacity>> {
Box::pin(self.active_provider.provider_capacities_for_input(input_name))
}
fn acquire<'a>(&'a self, input_name: &'a Arc<str>, priority: i8) -> BoxFuture<'a, Option<ProviderHandle>> {
Box::pin(self.active_provider.acquire_connection_for_download(input_name, priority))
}
fn release(&self, handle: Option<ProviderHandle>) -> BoxFuture<'_, ()> {
Box::pin(self.connection_manager.release_provider_handle(handle))
}
fn capacity_changed(&self) -> Arc<Notify> { self.connection_manager.capacity_notified() }
}
@@ -1489,6 +1489,10 @@ mod tests {
let (manual_update_sender, _) = mpsc::channel::<crate::api::model::ManualPlaylistUpdateRequest>(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::<crate::api::model::ManualPlaylistUpdateRequest>(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(),
+1
View File
@@ -1,3 +1,4 @@
pub mod recording_capacity;
pub mod recording_catalog_access;
pub mod recording_conflict;
pub mod recording_ctx;
@@ -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<str>, 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<str>) -> BoxFuture<'a, Vec<ProviderCapacity>>;
/// 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<str>, priority: i8) -> BoxFuture<'a, Option<ProviderHandle>>;
/// 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<ProviderHandle>) -> BoxFuture<'_, ()>;
/// Fires when a connection is freed anywhere, so waiters can re-check.
fn capacity_changed(&self) -> Arc<Notify>;
}
#[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<Vec<ProviderCapacity>>,
/// `None` means "no slot free", which is what makes a worker wait.
grants: Mutex<Vec<Option<ProviderHandle>>>,
notify: Arc<Notify>,
pub acquires: AtomicUsize,
pub releases: AtomicUsize,
}
impl StubCapacity {
fn with_capacity(in_use: usize, limit: usize) -> Arc<Self> {
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> { Self::with_capacity(0, 4) }
/// A provider that is full, so every request has to wait.
pub fn full() -> Arc<Self> { 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<str>) -> BoxFuture<'a, Vec<ProviderCapacity>> {
Box::pin(async move { self.capacities.lock().expect("capacities").clone() })
}
fn acquire<'a>(&'a self, _input_name: &'a Arc<str>, _priority: i8) -> BoxFuture<'a, Option<ProviderHandle>> {
Box::pin(async move {
self.acquires.fetch_add(1, Ordering::SeqCst);
self.grants.lock().expect("grants").pop().flatten()
})
}
fn release(&self, _handle: Option<ProviderHandle>) -> BoxFuture<'_, ()> {
Box::pin(async move {
self.releases.fetch_add(1, Ordering::SeqCst);
})
}
fn capacity_changed(&self) -> Arc<Notify> { 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<str> = 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);
}
}
+6 -7
View File
@@ -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<EventManager>,
/// Shared HTTP client, swapped when the proxy configuration changes.
pub http_client: Arc<ArcSwap<Client>>,
/// Provider capacity, used to acquire a connection slot before a
/// recording starts and to release it when the task ends.
pub active_provider: Arc<ActiveProviderManager>,
/// Connection registry; its capacity signal wakes waiting recordings.
pub connection_manager: Arc<ConnectionManager>,
/// 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<dyn RecordingCapacityPort>,
}
impl RecordingCtx {
+22 -29
View File
@@ -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<RecordingQueue>,
event_manager: &Arc<EventManager>,
active_provider: &Arc<ActiveProviderManager>,
connection_manager: &Arc<ConnectionManager>,
capacity: &Arc<dyn RecordingCapacityPort>,
) -> 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<RecordingQueue>,
event_manager: Arc<EventManager>,
active_provider: Arc<ActiveProviderManager>,
connection_manager: Arc<ConnectionManager>,
capacity: Arc<dyn RecordingCapacityPort>,
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;
}