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.
This commit is contained in:
DarkBreakpoint
2026-09-01 10:26:20 -05:00
parent af679d637d
commit 044fe1a70f
2 changed files with 87 additions and 38 deletions
+81 -4
View File
@@ -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);
@@ -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)));
}