From 044fe1a70f5fc2ac6f89c1bd1fdbdfccc8d5e286 Mon Sep 17 00:00:00 2001 From: DarkBreakpoint Date: Tue, 1 Sep 2026 10:26:20 -0500 Subject: [PATCH] fix(dvr): every path into the active slot attaches instead of re-running Making promotion materialization-aware only covered `promote_next_download`. Five other sites filled the active slot with `queue.remove(0)` directly, and one of them is the path that actually matters: when a transfer finishes, `finish_active_and_promote` took the next entry blindly. Two users queueing the same film hit exactly that -- the first completes, the second is promoted, and it downloads the same bytes over the file that was just written. The other four are the same shape: retry requeue, capacity-wait requeue, retry-limit-exceeded, and cancelling a paused active recording. Promotion is now one routine, `promote_from_queue`, and every site calls it. Filling the active slot without consulting the attach rule is no longer something a caller can do by accident. The two requeue sites gain something from this beyond consistency: an entry retrying a transfer whose file another entry finished in the meantime now attaches to it rather than retrying a download of bytes that are already on disk. `PromotionDecision::Wait` is documented as unreachable rather than left looking load-bearing. With a single active slot, promotion only runs once that slot is empty, so nothing can be in flight for the same media. It is kept because it is the correct answer if that ever changes. --- backend/dvr/src/recording/recording_queue.rs | 85 ++++++++++++++++++- .../dvr/src/recording/recording_transfer.rs | 40 ++------- 2 files changed, 87 insertions(+), 38 deletions(-) diff --git a/backend/dvr/src/recording/recording_queue.rs b/backend/dvr/src/recording/recording_queue.rs index ac267a7a2..181a77132 100644 --- a/backend/dvr/src/recording/recording_queue.rs +++ b/backend/dvr/src/recording/recording_queue.rs @@ -416,7 +416,9 @@ pub enum PromotionDecision { Execute, /// A completed entry at this index in `finished` already holds it. AttachTo(usize), - /// Another entry is producing it right now. + /// Another entry is producing it right now. Unreachable while there is a + /// single active slot, since promotion only runs once it is empty; kept so + /// the rule stays correct if that ever stops being true. Wait, } @@ -465,6 +467,35 @@ pub fn promotion_decision(candidate: &PersistedRecordingQueue, task: &PersistedR PromotionDecision::Execute } +/// Move the next runnable entry into the active slot, attaching any entry whose +/// file another entry already produced. +/// +/// Every path that fills the active slot goes through here. Taking the head of +/// the queue directly would re-download a file that was just completed by the +/// entry ahead of it. +pub fn promote_from_queue(candidate: &mut PersistedRecordingQueue) -> Option<(String, String)> { + let mut index = 0; + while index < candidate.queue.len() { + match promotion_decision(candidate, &candidate.queue[index]) { + PromotionDecision::Execute => { + let next = candidate.queue.remove(index); + let promoted = (next.uuid.clone(), next.filename.clone()); + candidate.active = Some(next); + return Some(promoted); + } + PromotionDecision::AttachTo(completed) => { + let source = candidate.finished[completed].clone(); + let mut attached = candidate.queue.remove(index); + attach_to_completed(&mut attached, &source); + candidate.finished.push(attached); + // The queue shrank; re-examine this position. + } + PromotionDecision::Wait => index += 1, + } + } + None +} + /// Adopt an already-produced file instead of transferring it again. /// /// Only the physical result is copied. The owner, visibility and quota belong @@ -1186,9 +1217,7 @@ impl RecordingQueue { cancelled.error.get_or_insert_with(|| "Cancelled by user".to_string()); cancelled.state = RecordingTaskState::Cancelled; candidate.finished.push(cancelled); - if !candidate.queue.is_empty() { - candidate.active = Some(candidate.queue.remove(0)); - } + promote_from_queue(candidate); } else if let Some(active) = candidate.active.as_mut() { active.state = RecordingTaskState::Cancelled; active.error = Some("Cancelled by user".to_string()); @@ -1507,6 +1536,54 @@ mod tests { assert_eq!(promotion_decision(&candidate, &queued), PromotionDecision::Execute); } + #[test] + fn the_entry_behind_a_finished_transfer_attaches_instead_of_downloading_again() { + // The live path: Alice's transfer completes and Bob is next in the + // queue for the same film. Taking the queue head blindly would start a + // second transfer over the file Alice just produced. + let mut candidate = PersistedRecordingQueue { + finished: vec![identified("alice", "film-42", RecordingTaskState::Completed)], + queue: vec![identified("bob", "film-42", RecordingTaskState::Queued)], + ..PersistedRecordingQueue::default() + }; + let promoted = promote_from_queue(&mut candidate); + assert!(promoted.is_none(), "nothing needs to run"); + assert!(candidate.active.is_none()); + assert!(candidate.queue.is_empty(), "bob left the queue"); + let bob = candidate.finished.iter().find(|task| task.uuid == "bob").expect("bob is filed"); + assert_eq!(bob.state, RecordingTaskState::Completed); + } + + #[test] + fn an_entry_for_other_media_behind_a_finished_transfer_still_runs() { + let mut candidate = PersistedRecordingQueue { + finished: vec![identified("alice", "film-42", RecordingTaskState::Completed)], + queue: vec![identified("bob", "film-99", RecordingTaskState::Queued)], + ..PersistedRecordingQueue::default() + }; + let promoted = promote_from_queue(&mut candidate); + assert_eq!(promoted.map(|(uuid, _)| uuid), Some("bob".to_string())); + assert_eq!(candidate.active.as_ref().map(|task| task.uuid.as_str()), Some("bob")); + } + + #[test] + fn attachable_entries_are_skipped_to_reach_one_that_must_run() { + let mut candidate = PersistedRecordingQueue { + finished: vec![identified("alice", "film-42", RecordingTaskState::Completed)], + queue: vec![ + identified("bob", "film-42", RecordingTaskState::Queued), + identified("carol", "film-42", RecordingTaskState::Queued), + identified("dave", "film-99", RecordingTaskState::Queued), + ], + ..PersistedRecordingQueue::default() + }; + let promoted = promote_from_queue(&mut candidate); + assert_eq!(promoted.map(|(uuid, _)| uuid), Some("dave".to_string())); + assert!(candidate.queue.is_empty(), "every entry was dispatched"); + // Bob and Carol were filed against Alice's file without running. + assert_eq!(candidate.finished.len(), 3); + } + #[test] fn attaching_adopts_the_file_but_keeps_the_entry_its_own() { let mut source = identified("a", "film-42", RecordingTaskState::Completed); diff --git a/backend/dvr/src/recording/recording_transfer.rs b/backend/dvr/src/recording/recording_transfer.rs index 5fd8abacf..c4f086834 100644 --- a/backend/dvr/src/recording/recording_transfer.rs +++ b/backend/dvr/src/recording/recording_transfer.rs @@ -919,7 +919,7 @@ async fn requeue_active_download_for_retry( download.next_retry_at = None; candidate.queue.insert(0, download); if promote { - candidate.active = Some(candidate.queue.remove(0)); + crate::recording::recording_queue::promote_from_queue(candidate); } Ok(Some(true)) }) @@ -948,7 +948,7 @@ async fn requeue_active_download_for_capacity_wait( download.next_retry_at = None; candidate.queue.insert(0, download); if promote { - candidate.active = Some(candidate.queue.remove(0)); + crate::recording::recording_queue::promote_from_queue(candidate); } Ok(Some(true)) }; @@ -967,29 +967,7 @@ async fn promote_next_download( if candidate.active.is_some() || candidate.queue.is_empty() { return Ok(None); } - // Walk the queue rather than taking the head unconditionally: an entry - // whose file another entry already produced must attach to it, and one - // whose file is being produced right now has to wait its turn. - let mut index = 0; - while index < candidate.queue.len() { - match crate::recording::recording_queue::promotion_decision(candidate, &candidate.queue[index]) { - crate::recording::recording_queue::PromotionDecision::Execute => { - let next = candidate.queue.remove(index); - let promoted = (next.uuid.clone(), next.filename.clone()); - candidate.active = Some(next); - return Ok(Some(promoted)); - } - crate::recording::recording_queue::PromotionDecision::AttachTo(completed) => { - let source = candidate.finished[completed].clone(); - let mut attached = candidate.queue.remove(index); - crate::recording::recording_queue::attach_to_completed(&mut attached, &source); - candidate.finished.push(attached); - // The queue shrank; re-examine this position. - } - crate::recording::recording_queue::PromotionDecision::Wait => index += 1, - } - } - Ok(None) + Ok(crate::recording::recording_queue::promote_from_queue(candidate)) }) .await } @@ -1012,9 +990,7 @@ where } let notification = finish(&mut active); candidate.finished.push(active); - if !candidate.queue.is_empty() { - candidate.active = Some(candidate.queue.remove(0)); - } + crate::recording::recording_queue::promote_from_queue(candidate); Ok(Some(notification)) }) .await @@ -1036,9 +1012,7 @@ async fn cancel_active_and_promote(download_queue: &RecordingQueue, uuid: &str) active.error.get_or_insert_with(|| "Cancelled by user".to_string()); active.state = RecordingTaskState::Cancelled; candidate.finished.push(active); - if !candidate.queue.is_empty() { - candidate.active = Some(candidate.queue.remove(0)); - } + crate::recording::recording_queue::promote_from_queue(candidate); Ok(Some(true)) }) .await? @@ -1074,9 +1048,7 @@ async fn prepare_active_retry( let notification = mark_recording_metadata_notification(&mut failed.recording, LifecycleEvent::Failed, Some(error)); candidate.finished.push(failed); - if !candidate.queue.is_empty() { - candidate.active = Some(candidate.queue.remove(0)); - } + crate::recording::recording_queue::promote_from_queue(candidate); return Ok(Some(RetryCommit::Failed(notification))); }