Jeez...
This commit is contained in:
parent
efcffdd072
commit
51d1a0b0c7
7 changed files with 183 additions and 115 deletions
|
|
@ -14,7 +14,9 @@ pub struct ChunkTracker {
|
||||||
//
|
//
|
||||||
// Initially this is the opposite of "have", until we start making requests.
|
// Initially this is the opposite of "have", until we start making requests.
|
||||||
// An in-flight request is not in "needed", and not in "have".
|
// An in-flight request is not in "needed", and not in "have".
|
||||||
needed_pieces: BF,
|
//
|
||||||
|
// needed initial value = selected & !have
|
||||||
|
queue_pieces: BF,
|
||||||
|
|
||||||
// This has a bit set per each chunk (block) that we have written to the output file.
|
// This has a bit set per each chunk (block) that we have written to the output file.
|
||||||
// It doesn't mean it's valid yet. Used to track how much is left in each piece.
|
// It doesn't mean it's valid yet. Used to track how much is left in each piece.
|
||||||
|
|
@ -23,18 +25,35 @@ pub struct ChunkTracker {
|
||||||
// These are the pieces that we actually have, fully checked and downloaded.
|
// These are the pieces that we actually have, fully checked and downloaded.
|
||||||
have: BF,
|
have: BF,
|
||||||
|
|
||||||
|
// The pieces that the user selected. This doesn't change unless update_only_files
|
||||||
|
// was called.
|
||||||
|
selected: BF,
|
||||||
|
|
||||||
lengths: Lengths,
|
lengths: Lengths,
|
||||||
|
|
||||||
// What pieces to download first.
|
// What pieces to download first.
|
||||||
priority_piece_ids: Vec<usize>,
|
priority_piece_ids: Vec<usize>,
|
||||||
|
|
||||||
total_selected_bytes: u64,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, PartialEq, Eq, Clone, Copy)]
|
#[derive(Debug, PartialEq, Eq, Clone, Copy)]
|
||||||
pub struct HaveNeeded {
|
pub struct HaveNeededSelected {
|
||||||
pub have_bytes: u64,
|
pub have_bytes: u64,
|
||||||
pub needed_bytes: u64,
|
pub needed_bytes: u64,
|
||||||
|
pub selected_bytes: u64,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl HaveNeededSelected {
|
||||||
|
pub const fn progress(&self) -> u64 {
|
||||||
|
self.selected_bytes - self.needed_bytes
|
||||||
|
}
|
||||||
|
|
||||||
|
pub const fn total(&self) -> u64 {
|
||||||
|
self.selected_bytes
|
||||||
|
}
|
||||||
|
|
||||||
|
pub const fn finished(&self) -> bool {
|
||||||
|
self.needed_bytes == 0
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Comput the have-status of chunks.
|
// Comput the have-status of chunks.
|
||||||
|
|
@ -68,6 +87,29 @@ fn compute_chunk_have_status(lengths: &Lengths, have_pieces: &BF) -> anyhow::Res
|
||||||
Ok(chunk_bf)
|
Ok(chunk_bf)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn compute_queued_pieces_unchecked(have_pieces: &BF, selected_pieces: &BF) -> BF {
|
||||||
|
// it's needed ONLY if it's selected and we don't have it.
|
||||||
|
use core::ops::BitAnd;
|
||||||
|
use core::ops::Not;
|
||||||
|
|
||||||
|
have_pieces.clone().not().bitand(selected_pieces)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn compute_queued_pieces(have_pieces: &BF, selected_pieces: &BF) -> anyhow::Result<BF> {
|
||||||
|
if have_pieces.len() != selected_pieces.len() {
|
||||||
|
anyhow::bail!(
|
||||||
|
"have_pieces.len() != selected_pieces.len(), {} != {}",
|
||||||
|
have_pieces.len(),
|
||||||
|
selected_pieces.len()
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(compute_queued_pieces_unchecked(
|
||||||
|
have_pieces,
|
||||||
|
selected_pieces,
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|
||||||
pub enum ChunkMarkingResult {
|
pub enum ChunkMarkingResult {
|
||||||
PreviouslyCompleted,
|
PreviouslyCompleted,
|
||||||
NotCompleted,
|
NotCompleted,
|
||||||
|
|
@ -76,15 +118,15 @@ pub enum ChunkMarkingResult {
|
||||||
|
|
||||||
impl ChunkTracker {
|
impl ChunkTracker {
|
||||||
pub fn new(
|
pub fn new(
|
||||||
// Needed pieces are the ones we need to download. NOTE: if all files are selected,
|
|
||||||
// this is the inverse of have_pieces. But if partial files are selected, we may need more/less
|
|
||||||
// than we have.
|
|
||||||
needed_pieces: BF,
|
|
||||||
// Have pieces are the ones we have already downloaded and verified.
|
// Have pieces are the ones we have already downloaded and verified.
|
||||||
have_pieces: BF,
|
have_pieces: BF,
|
||||||
|
// Selected pieces are the ones the user has selected
|
||||||
|
selected_pieces: BF,
|
||||||
lengths: Lengths,
|
lengths: Lengths,
|
||||||
total_selected_bytes: u64,
|
|
||||||
) -> anyhow::Result<Self> {
|
) -> anyhow::Result<Self> {
|
||||||
|
let needed_pieces = compute_queued_pieces(&have_pieces, &selected_pieces)
|
||||||
|
.context("error computing needed pieces")?;
|
||||||
|
|
||||||
// TODO: ideally this needs to be a list based on needed files, e.g.
|
// TODO: ideally this needs to be a list based on needed files, e.g.
|
||||||
// last needed piece for each file. But let's keep simple for now.
|
// last needed piece for each file. But let's keep simple for now.
|
||||||
|
|
||||||
|
|
@ -103,18 +145,14 @@ impl ChunkTracker {
|
||||||
Ok(Self {
|
Ok(Self {
|
||||||
chunk_status: compute_chunk_have_status(&lengths, &have_pieces)
|
chunk_status: compute_chunk_have_status(&lengths, &have_pieces)
|
||||||
.context("error computing chunk status")?,
|
.context("error computing chunk status")?,
|
||||||
needed_pieces,
|
queue_pieces: needed_pieces,
|
||||||
|
selected: selected_pieces,
|
||||||
lengths,
|
lengths,
|
||||||
have: have_pieces,
|
have: have_pieces,
|
||||||
priority_piece_ids,
|
priority_piece_ids,
|
||||||
total_selected_bytes,
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn get_total_selected_bytes(&self) -> u64 {
|
|
||||||
self.total_selected_bytes
|
|
||||||
}
|
|
||||||
|
|
||||||
pub fn get_lengths(&self) -> &Lengths {
|
pub fn get_lengths(&self) -> &Lengths {
|
||||||
&self.lengths
|
&self.lengths
|
||||||
}
|
}
|
||||||
|
|
@ -123,7 +161,7 @@ impl ChunkTracker {
|
||||||
&self.have
|
&self.have
|
||||||
}
|
}
|
||||||
pub fn reserve_needed_piece(&mut self, index: ValidPieceIndex) {
|
pub fn reserve_needed_piece(&mut self, index: ValidPieceIndex) {
|
||||||
self.needed_pieces.set(index.get() as usize, false)
|
self.queue_pieces.set(index.get() as usize, false)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn calc_have_bytes(&self) -> u64 {
|
pub fn calc_have_bytes(&self) -> u64 {
|
||||||
|
|
@ -137,22 +175,28 @@ impl ChunkTracker {
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn calc_needed_bytes(&self) -> u64 {
|
pub fn calc_needed_bytes(&self) -> u64 {
|
||||||
self.needed_pieces
|
self.have
|
||||||
.iter_ones()
|
.iter()
|
||||||
.filter_map(|piece_id| {
|
.zip(self.selected.iter())
|
||||||
let piece_id = self.lengths.validate_piece_index(piece_id as u32)?;
|
.enumerate()
|
||||||
Some(self.lengths.piece_length(piece_id) as u64)
|
.filter_map(|(piece_id, (have, selected))| {
|
||||||
|
if *selected && !*have {
|
||||||
|
let piece_id = self.lengths.validate_piece_index(piece_id as u32)?;
|
||||||
|
Some(self.lengths.piece_length(piece_id) as u64)
|
||||||
|
} else {
|
||||||
|
None
|
||||||
|
}
|
||||||
})
|
})
|
||||||
.sum()
|
.sum()
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn iter_needed_pieces(&self) -> impl Iterator<Item = usize> + '_ {
|
pub fn iter_queued_pieces(&self) -> impl Iterator<Item = usize> + '_ {
|
||||||
self.priority_piece_ids
|
self.priority_piece_ids
|
||||||
.iter()
|
.iter()
|
||||||
.copied()
|
.copied()
|
||||||
.filter(move |piece_id| self.needed_pieces[*piece_id])
|
.filter(move |piece_id| self.queue_pieces[*piece_id])
|
||||||
.chain(
|
.chain(
|
||||||
self.needed_pieces
|
self.queue_pieces
|
||||||
.iter_ones()
|
.iter_ones()
|
||||||
.filter(move |id| !self.priority_piece_ids.contains(id)),
|
.filter(move |id| !self.priority_piece_ids.contains(id)),
|
||||||
)
|
)
|
||||||
|
|
@ -172,7 +216,7 @@ impl ChunkTracker {
|
||||||
// This will trigger the requesters to re-check each chunk in this piece.
|
// This will trigger the requesters to re-check each chunk in this piece.
|
||||||
let chunk_range = self.lengths.chunk_range(index);
|
let chunk_range = self.lengths.chunk_range(index);
|
||||||
if !self.chunk_status.get(chunk_range)?.all() {
|
if !self.chunk_status.get(chunk_range)?.all() {
|
||||||
self.needed_pieces.set(index.get() as usize, true);
|
self.queue_pieces.set(index.get() as usize, true);
|
||||||
}
|
}
|
||||||
Some(true)
|
Some(true)
|
||||||
}
|
}
|
||||||
|
|
@ -187,7 +231,7 @@ impl ChunkTracker {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
debug!("remarking piece={} as broken", index);
|
debug!("remarking piece={} as broken", index);
|
||||||
self.needed_pieces.set(index.get() as usize, true);
|
self.queue_pieces.set(index.get() as usize, true);
|
||||||
if let Some(s) = self.chunk_status.get_mut(self.lengths.chunk_range(index)) {
|
if let Some(s) = self.chunk_status.get_mut(self.lengths.chunk_range(index)) {
|
||||||
s.fill(false);
|
s.fill(false);
|
||||||
}
|
}
|
||||||
|
|
@ -244,14 +288,15 @@ impl ChunkTracker {
|
||||||
file_lengths_iterator: impl IntoIterator<Item = u64>,
|
file_lengths_iterator: impl IntoIterator<Item = u64>,
|
||||||
// TODO: maybe make this a BF
|
// TODO: maybe make this a BF
|
||||||
new_only_files: &HashSet<usize>,
|
new_only_files: &HashSet<usize>,
|
||||||
) -> anyhow::Result<HaveNeeded> {
|
) -> anyhow::Result<HaveNeededSelected> {
|
||||||
let mut piece_it = self.lengths.iter_piece_infos();
|
let mut piece_it = self.lengths.iter_piece_infos();
|
||||||
let mut current_piece = piece_it
|
let mut current_piece = piece_it
|
||||||
.next()
|
.next()
|
||||||
.context("bug: iter_piece_infos() returned empty iterator")?;
|
.context("bug: iter_piece_infos() returned empty iterator")?;
|
||||||
let mut current_piece_needed = false;
|
let mut current_piece_selected = false;
|
||||||
let mut current_piece_remaining = current_piece.len;
|
let mut current_piece_remaining = current_piece.len;
|
||||||
let mut have_bytes = 0u64;
|
let mut have_bytes = 0u64;
|
||||||
|
let mut selected_bytes = 0u64;
|
||||||
let mut needed_bytes = 0u64;
|
let mut needed_bytes = 0u64;
|
||||||
|
|
||||||
for (idx, len) in file_lengths_iterator.into_iter().enumerate() {
|
for (idx, len) in file_lengths_iterator.into_iter().enumerate() {
|
||||||
|
|
@ -260,31 +305,38 @@ impl ChunkTracker {
|
||||||
let mut remaining_file_len = len;
|
let mut remaining_file_len = len;
|
||||||
|
|
||||||
while remaining_file_len > 0 {
|
while remaining_file_len > 0 {
|
||||||
current_piece_needed |= len > 0 && file_required;
|
current_piece_selected |= len > 0 && file_required;
|
||||||
let shift = std::cmp::min(current_piece_remaining as u64, remaining_file_len);
|
let shift = std::cmp::min(current_piece_remaining as u64, remaining_file_len);
|
||||||
assert!(shift > 0);
|
assert!(shift > 0);
|
||||||
remaining_file_len -= shift;
|
remaining_file_len -= shift;
|
||||||
current_piece_remaining -= shift as u32;
|
current_piece_remaining -= shift as u32;
|
||||||
|
|
||||||
dbg!(
|
// dbg!(
|
||||||
idx,
|
// idx,
|
||||||
shift,
|
// shift,
|
||||||
remaining_file_len,
|
// remaining_file_len,
|
||||||
current_piece_remaining,
|
// current_piece_remaining,
|
||||||
current_piece_needed,
|
// current_piece_needed,
|
||||||
file_required,
|
// file_required,
|
||||||
current_piece
|
// current_piece
|
||||||
);
|
// );
|
||||||
|
|
||||||
if current_piece_remaining == 0 {
|
if current_piece_remaining == 0 {
|
||||||
let current_piece_have = self.have[current_piece.piece_index.get() as usize];
|
let current_piece_have = self.have[current_piece.piece_index.get() as usize];
|
||||||
if current_piece_have {
|
if current_piece_have {
|
||||||
have_bytes += current_piece.len as u64;
|
have_bytes += current_piece.len as u64;
|
||||||
}
|
}
|
||||||
if current_piece_needed {
|
if current_piece_selected {
|
||||||
|
selected_bytes += current_piece.len as u64;
|
||||||
|
}
|
||||||
|
if current_piece_selected && !current_piece_have {
|
||||||
needed_bytes += current_piece.len as u64;
|
needed_bytes += current_piece.len as u64;
|
||||||
}
|
}
|
||||||
match (current_piece_needed, current_piece_have) {
|
self.selected.set(
|
||||||
|
current_piece.piece_index.get() as usize,
|
||||||
|
current_piece_selected,
|
||||||
|
);
|
||||||
|
match (current_piece_selected, current_piece_have) {
|
||||||
(true, true) => {}
|
(true, true) => {}
|
||||||
(true, false) => {
|
(true, false) => {
|
||||||
dbg!(self.mark_piece_broken_if_not_have(current_piece.piece_index))
|
dbg!(self.mark_piece_broken_if_not_have(current_piece.piece_index))
|
||||||
|
|
@ -293,7 +345,7 @@ impl ChunkTracker {
|
||||||
(false, false) => {
|
(false, false) => {
|
||||||
// don't need the piece, and don't have it - cancel downloading it
|
// don't need the piece, and don't have it - cancel downloading it
|
||||||
dbg!(self
|
dbg!(self
|
||||||
.needed_pieces
|
.queue_pieces
|
||||||
.set(current_piece.piece_index.get() as usize, false));
|
.set(current_piece.piece_index.get() as usize, false));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -302,16 +354,17 @@ impl ChunkTracker {
|
||||||
current_piece = piece_it.next().context(
|
current_piece = piece_it.next().context(
|
||||||
"bug: iter_piece_infos() pieces ended earlier than expected",
|
"bug: iter_piece_infos() pieces ended earlier than expected",
|
||||||
)?;
|
)?;
|
||||||
current_piece_needed = false;
|
current_piece_selected = false;
|
||||||
current_piece_remaining = current_piece.len;
|
current_piece_remaining = current_piece.len;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
Ok(HaveNeeded {
|
Ok(HaveNeededSelected {
|
||||||
have_bytes,
|
have_bytes,
|
||||||
needed_bytes,
|
needed_bytes,
|
||||||
|
selected_bytes,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -322,7 +375,7 @@ mod tests {
|
||||||
|
|
||||||
use librqbit_core::{constants::CHUNK_SIZE, lengths::Lengths};
|
use librqbit_core::{constants::CHUNK_SIZE, lengths::Lengths};
|
||||||
|
|
||||||
use crate::{chunk_tracker::HaveNeeded, type_aliases::BF};
|
use crate::{chunk_tracker::HaveNeededSelected, type_aliases::BF};
|
||||||
|
|
||||||
use super::{compute_chunk_have_status, ChunkTracker};
|
use super::{compute_chunk_have_status, ChunkTracker};
|
||||||
|
|
||||||
|
|
@ -444,103 +497,104 @@ mod tests {
|
||||||
let initial_needed = BF::from_boxed_slice(vec![u8::MAX; bf_len].into_boxed_slice());
|
let initial_needed = BF::from_boxed_slice(vec![u8::MAX; bf_len].into_boxed_slice());
|
||||||
|
|
||||||
// Initially, we need all files and all pieces.
|
// Initially, we need all files and all pieces.
|
||||||
let mut ct = ChunkTracker::new(
|
let mut ct = ChunkTracker::new(initial_needed.clone(), initial_have.clone(), l).unwrap();
|
||||||
initial_needed.clone(),
|
|
||||||
initial_have.clone(),
|
|
||||||
l,
|
|
||||||
l.total_length(),
|
|
||||||
)
|
|
||||||
.unwrap();
|
|
||||||
|
|
||||||
// Select all file, no changes.
|
// Select all file, no changes.
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
ct.update_only_files(all_files.into_iter(), &HashSet::from_iter([0, 1, 2, 3]))
|
ct.update_only_files(all_files.into_iter(), &HashSet::from_iter([0, 1, 2, 3]))
|
||||||
.unwrap(),
|
.unwrap(),
|
||||||
HaveNeeded {
|
HaveNeededSelected {
|
||||||
have_bytes: 0,
|
have_bytes: 0,
|
||||||
needed_bytes: total_len
|
selected_bytes: total_len,
|
||||||
|
needed_bytes: total_len,
|
||||||
}
|
}
|
||||||
);
|
);
|
||||||
assert_eq!(ct.have, initial_have);
|
assert_eq!(ct.have, initial_have);
|
||||||
assert_eq!(ct.needed_pieces, initial_needed);
|
assert_eq!(ct.queue_pieces, initial_needed);
|
||||||
|
|
||||||
// Select only the first file.
|
// Select only the first file.
|
||||||
println!("Select only the first file.");
|
println!("Select only the first file.");
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
ct.update_only_files(all_files, &HashSet::from_iter([0]))
|
ct.update_only_files(all_files, &HashSet::from_iter([0]))
|
||||||
.unwrap(),
|
.unwrap(),
|
||||||
HaveNeeded {
|
HaveNeededSelected {
|
||||||
have_bytes: 0,
|
have_bytes: 0,
|
||||||
|
selected_bytes: all_files[0],
|
||||||
needed_bytes: all_files[0],
|
needed_bytes: all_files[0],
|
||||||
}
|
}
|
||||||
);
|
);
|
||||||
assert_eq!(ct.needed_pieces[0], true);
|
assert_eq!(ct.queue_pieces[0], true);
|
||||||
assert_eq!(ct.needed_pieces[1], false);
|
assert_eq!(ct.queue_pieces[1], false);
|
||||||
assert_eq!(ct.needed_pieces[2], false);
|
assert_eq!(ct.queue_pieces[2], false);
|
||||||
|
|
||||||
// Select only the second file.
|
// Select only the second file.
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
ct.update_only_files(all_files, &HashSet::from_iter([1]))
|
ct.update_only_files(all_files, &HashSet::from_iter([1]))
|
||||||
.unwrap(),
|
.unwrap(),
|
||||||
HaveNeeded {
|
HaveNeededSelected {
|
||||||
have_bytes: 0,
|
have_bytes: 0,
|
||||||
|
selected_bytes: piece_len as u64,
|
||||||
needed_bytes: piece_len as u64,
|
needed_bytes: piece_len as u64,
|
||||||
}
|
}
|
||||||
);
|
);
|
||||||
assert_eq!(ct.needed_pieces[0], false);
|
assert_eq!(ct.queue_pieces[0], false);
|
||||||
assert_eq!(ct.needed_pieces[1], true);
|
assert_eq!(ct.queue_pieces[1], true);
|
||||||
assert_eq!(ct.needed_pieces[2], false);
|
assert_eq!(ct.queue_pieces[2], false);
|
||||||
|
|
||||||
// Select only the third file (zero sized one!).
|
// Select only the third file (zero sized one!).
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
ct.update_only_files(all_files, &HashSet::from_iter([2]))
|
ct.update_only_files(all_files, &HashSet::from_iter([2]))
|
||||||
.unwrap(),
|
.unwrap(),
|
||||||
HaveNeeded {
|
HaveNeededSelected {
|
||||||
have_bytes: 0,
|
have_bytes: 0,
|
||||||
|
selected_bytes: 0,
|
||||||
needed_bytes: 0,
|
needed_bytes: 0,
|
||||||
}
|
}
|
||||||
);
|
);
|
||||||
assert_eq!(ct.needed_pieces[0], false);
|
assert_eq!(ct.queue_pieces[0], false);
|
||||||
assert_eq!(ct.needed_pieces[1], false);
|
assert_eq!(ct.queue_pieces[1], false);
|
||||||
assert_eq!(ct.needed_pieces[2], false);
|
assert_eq!(ct.queue_pieces[2], false);
|
||||||
|
|
||||||
// Select only the fourth file.
|
// Select only the fourth file.
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
ct.update_only_files(all_files, &HashSet::from_iter([3]))
|
ct.update_only_files(all_files, &HashSet::from_iter([3]))
|
||||||
.unwrap(),
|
.unwrap(),
|
||||||
HaveNeeded {
|
HaveNeededSelected {
|
||||||
have_bytes: 0,
|
have_bytes: 0,
|
||||||
|
selected_bytes: (piece_len + 1) as u64,
|
||||||
needed_bytes: (piece_len + 1) as u64,
|
needed_bytes: (piece_len + 1) as u64,
|
||||||
}
|
}
|
||||||
);
|
);
|
||||||
assert_eq!(ct.needed_pieces[0], false);
|
assert_eq!(ct.queue_pieces[0], false);
|
||||||
assert_eq!(ct.needed_pieces[1], true);
|
assert_eq!(ct.queue_pieces[1], true);
|
||||||
assert_eq!(ct.needed_pieces[2], true);
|
assert_eq!(ct.queue_pieces[2], true);
|
||||||
|
|
||||||
// Select first and last file
|
// Select first and last file
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
ct.update_only_files(all_files.clone(), &HashSet::from_iter([0, 3]))
|
ct.update_only_files(all_files.clone(), &HashSet::from_iter([0, 3]))
|
||||||
.unwrap(),
|
.unwrap(),
|
||||||
HaveNeeded {
|
HaveNeededSelected {
|
||||||
have_bytes: 0,
|
have_bytes: 0,
|
||||||
|
selected_bytes: all_files[0] + all_files[3] + 1,
|
||||||
needed_bytes: all_files[0] + all_files[3] + 1,
|
needed_bytes: all_files[0] + all_files[3] + 1,
|
||||||
}
|
}
|
||||||
);
|
);
|
||||||
assert_eq!(ct.needed_pieces[0], true);
|
assert_eq!(ct.queue_pieces[0], true);
|
||||||
assert_eq!(ct.needed_pieces[1], true);
|
assert_eq!(ct.queue_pieces[1], true);
|
||||||
assert_eq!(ct.needed_pieces[2], true);
|
assert_eq!(ct.queue_pieces[2], true);
|
||||||
|
|
||||||
// Select all files
|
// Select all files
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
ct.update_only_files(all_files.clone(), &HashSet::from_iter([0, 1, 2, 3]))
|
ct.update_only_files(all_files.clone(), &HashSet::from_iter([0, 1, 2, 3]))
|
||||||
.unwrap(),
|
.unwrap(),
|
||||||
HaveNeeded {
|
HaveNeededSelected {
|
||||||
have_bytes: 0,
|
have_bytes: 0,
|
||||||
needed_bytes: total_len,
|
selected_bytes: total_len,
|
||||||
|
needed_bytes: total_len
|
||||||
}
|
}
|
||||||
);
|
);
|
||||||
assert_eq!(ct.needed_pieces[0], true);
|
assert_eq!(ct.queue_pieces[0], true);
|
||||||
assert_eq!(ct.needed_pieces[1], true);
|
assert_eq!(ct.queue_pieces[1], true);
|
||||||
assert_eq!(ct.needed_pieces[2], true);
|
assert_eq!(ct.queue_pieces[2], true);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -22,18 +22,26 @@ use tracing::{debug, trace, warn};
|
||||||
use crate::type_aliases::{PeerHandle, BF};
|
use crate::type_aliases::{PeerHandle, BF};
|
||||||
|
|
||||||
pub(crate) struct InitialCheckResults {
|
pub(crate) struct InitialCheckResults {
|
||||||
// The pieces that we need to download.
|
// A piece as flags based on these dimensions:
|
||||||
pub needed_pieces: BF,
|
// - if the asked for it or not (only_files)
|
||||||
|
// - if we have it downloaded and verified
|
||||||
|
// - if we need to queue it for downloading
|
||||||
|
// this one depends if we queued it already or not.
|
||||||
|
|
||||||
// The pieces we have downloaded.
|
// The pieces we have downloaded.
|
||||||
pub have_pieces: BF,
|
pub have_pieces: BF,
|
||||||
|
// The pieces that the user selected to download.
|
||||||
|
pub selected_pieces: BF,
|
||||||
|
|
||||||
// How many bytes we have. This can be MORE than "total_selected_bytes",
|
// How many bytes we have. This can be MORE than "total_selected_bytes",
|
||||||
// if we downloaded some pieces, and later the "only_files" was changed.
|
// if we downloaded some pieces, and later the "only_files" was changed.
|
||||||
pub have_bytes: u64,
|
pub have_bytes: u64,
|
||||||
// How many bytes we need to download.
|
// How many bytes we need to download.
|
||||||
pub needed_bytes: u64,
|
pub needed_bytes: u64,
|
||||||
|
|
||||||
// How many bytes are in selected pieces.
|
// How many bytes are in selected pieces.
|
||||||
// If all selected, this must be equal to total torrent length.
|
// If all selected, this must be equal to total torrent length.
|
||||||
pub total_selected_bytes: u64,
|
pub selected_bytes: u64,
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn update_hash_from_file<Sha1: ISha1>(
|
pub fn update_hash_from_file<Sha1: ISha1>(
|
||||||
|
|
@ -82,8 +90,8 @@ impl<'a> FileOps<'a> {
|
||||||
) -> anyhow::Result<InitialCheckResults> {
|
) -> anyhow::Result<InitialCheckResults> {
|
||||||
let mut needed_pieces =
|
let mut needed_pieces =
|
||||||
BF::from_boxed_slice(vec![0u8; self.lengths.piece_bitfield_bytes()].into());
|
BF::from_boxed_slice(vec![0u8; self.lengths.piece_bitfield_bytes()].into());
|
||||||
let mut have_pieces =
|
let mut have_pieces = needed_pieces.clone();
|
||||||
BF::from_boxed_slice(vec![0u8; self.lengths.piece_bitfield_bytes()].into());
|
let mut selected_pieces = needed_pieces.clone();
|
||||||
|
|
||||||
let mut have_bytes = 0u64;
|
let mut have_bytes = 0u64;
|
||||||
let mut needed_bytes = 0u64;
|
let mut needed_bytes = 0u64;
|
||||||
|
|
@ -139,7 +147,7 @@ impl<'a> FileOps<'a> {
|
||||||
let mut computed_hash = Sha1::new();
|
let mut computed_hash = Sha1::new();
|
||||||
let mut piece_remaining = piece_info.len as usize;
|
let mut piece_remaining = piece_info.len as usize;
|
||||||
let mut some_files_broken = false;
|
let mut some_files_broken = false;
|
||||||
let mut at_least_one_file_required = current_file.full_file_required;
|
let mut piece_selected = current_file.full_file_required;
|
||||||
progress.fetch_add(piece_info.len as u64, Ordering::Relaxed);
|
progress.fetch_add(piece_info.len as u64, Ordering::Relaxed);
|
||||||
|
|
||||||
while piece_remaining > 0 {
|
while piece_remaining > 0 {
|
||||||
|
|
@ -152,7 +160,7 @@ impl<'a> FileOps<'a> {
|
||||||
.next()
|
.next()
|
||||||
.ok_or_else(|| anyhow::anyhow!("broken torrent metadata"))?;
|
.ok_or_else(|| anyhow::anyhow!("broken torrent metadata"))?;
|
||||||
|
|
||||||
at_least_one_file_required |= current_file.full_file_required;
|
piece_selected |= current_file.full_file_required;
|
||||||
|
|
||||||
to_read_in_file =
|
to_read_in_file =
|
||||||
std::cmp::min(current_file.remaining(), piece_remaining as u64) as usize;
|
std::cmp::min(current_file.remaining(), piece_remaining as u64) as usize;
|
||||||
|
|
@ -186,18 +194,18 @@ impl<'a> FileOps<'a> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if at_least_one_file_required {
|
if piece_selected {
|
||||||
total_selected_bytes += piece_info.len as u64;
|
total_selected_bytes += piece_info.len as u64;
|
||||||
|
selected_pieces.set(piece_info.piece_index.get() as usize, true);
|
||||||
}
|
}
|
||||||
|
|
||||||
if at_least_one_file_required && some_files_broken {
|
if piece_selected && some_files_broken {
|
||||||
trace!(
|
trace!(
|
||||||
"piece {} had errors, marking as needed",
|
"piece {} had errors, marking as needed",
|
||||||
piece_info.piece_index
|
piece_info.piece_index
|
||||||
);
|
);
|
||||||
|
|
||||||
needed_bytes += piece_info.len as u64;
|
needed_bytes += piece_info.len as u64;
|
||||||
needed_pieces.set(piece_info.piece_index.get() as usize, true);
|
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -212,7 +220,7 @@ impl<'a> FileOps<'a> {
|
||||||
);
|
);
|
||||||
have_bytes += piece_info.len as u64;
|
have_bytes += piece_info.len as u64;
|
||||||
have_pieces.set(piece_info.piece_index.get() as usize, true);
|
have_pieces.set(piece_info.piece_index.get() as usize, true);
|
||||||
} else if at_least_one_file_required {
|
} else if piece_selected {
|
||||||
trace!(
|
trace!(
|
||||||
"piece {} hash does not match, marking as needed",
|
"piece {} hash does not match, marking as needed",
|
||||||
piece_info.piece_index
|
piece_info.piece_index
|
||||||
|
|
@ -228,11 +236,11 @@ impl<'a> FileOps<'a> {
|
||||||
}
|
}
|
||||||
|
|
||||||
Ok(InitialCheckResults {
|
Ok(InitialCheckResults {
|
||||||
needed_pieces,
|
|
||||||
have_pieces,
|
have_pieces,
|
||||||
|
selected_pieces,
|
||||||
have_bytes,
|
have_bytes,
|
||||||
needed_bytes,
|
needed_bytes,
|
||||||
total_selected_bytes,
|
selected_bytes: total_selected_bytes,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -253,7 +253,7 @@ async fn test_e2e() {
|
||||||
.with_state(|s| match s {
|
.with_state(|s| match s {
|
||||||
crate::ManagedTorrentState::Initializing(_) => Ok(false),
|
crate::ManagedTorrentState::Initializing(_) => Ok(false),
|
||||||
crate::ManagedTorrentState::Paused(p) => {
|
crate::ManagedTorrentState::Paused(p) => {
|
||||||
assert_eq!(p.needed_bytes, 0);
|
assert_eq!(p.hns.needed_bytes, 0);
|
||||||
Ok(true)
|
Ok(true)
|
||||||
}
|
}
|
||||||
_ => bail!("bugged state"),
|
_ => bail!("bugged state"),
|
||||||
|
|
|
||||||
|
|
@ -11,7 +11,10 @@ use parking_lot::Mutex;
|
||||||
use size_format::SizeFormatterBinary as SF;
|
use size_format::SizeFormatterBinary as SF;
|
||||||
use tracing::{debug, info, warn};
|
use tracing::{debug, info, warn};
|
||||||
|
|
||||||
use crate::{chunk_tracker::ChunkTracker, file_ops::FileOps};
|
use crate::{
|
||||||
|
chunk_tracker::{ChunkTracker, HaveNeededSelected},
|
||||||
|
file_ops::FileOps,
|
||||||
|
};
|
||||||
|
|
||||||
use super::{paused::TorrentStatePaused, ManagedTorrentInfo};
|
use super::{paused::TorrentStatePaused, ManagedTorrentInfo};
|
||||||
|
|
||||||
|
|
@ -88,7 +91,7 @@ impl TorrentStateInitializing {
|
||||||
"Initial check results: have {}, needed {}, total selected {}",
|
"Initial check results: have {}, needed {}, total selected {}",
|
||||||
SF::new(initial_check_results.have_bytes),
|
SF::new(initial_check_results.have_bytes),
|
||||||
SF::new(initial_check_results.needed_bytes),
|
SF::new(initial_check_results.needed_bytes),
|
||||||
SF::new(initial_check_results.total_selected_bytes)
|
SF::new(initial_check_results.selected_bytes)
|
||||||
);
|
);
|
||||||
|
|
||||||
self.meta.spawner.spawn_block_in_place(|| {
|
self.meta.spawner.spawn_block_in_place(|| {
|
||||||
|
|
@ -123,10 +126,9 @@ impl TorrentStateInitializing {
|
||||||
});
|
});
|
||||||
|
|
||||||
let chunk_tracker = ChunkTracker::new(
|
let chunk_tracker = ChunkTracker::new(
|
||||||
initial_check_results.needed_pieces,
|
|
||||||
initial_check_results.have_pieces,
|
initial_check_results.have_pieces,
|
||||||
|
initial_check_results.selected_pieces,
|
||||||
self.meta.lengths,
|
self.meta.lengths,
|
||||||
initial_check_results.total_selected_bytes,
|
|
||||||
)
|
)
|
||||||
.context("error creating chunk tracker")?;
|
.context("error creating chunk tracker")?;
|
||||||
|
|
||||||
|
|
@ -135,8 +137,11 @@ impl TorrentStateInitializing {
|
||||||
files,
|
files,
|
||||||
filenames,
|
filenames,
|
||||||
chunk_tracker,
|
chunk_tracker,
|
||||||
have_bytes: initial_check_results.have_bytes,
|
hns: HaveNeededSelected {
|
||||||
needed_bytes: initial_check_results.needed_bytes,
|
have_bytes: initial_check_results.have_bytes,
|
||||||
|
needed_bytes: initial_check_results.needed_bytes,
|
||||||
|
selected_bytes: initial_check_results.selected_bytes,
|
||||||
|
},
|
||||||
};
|
};
|
||||||
Ok(paused)
|
Ok(paused)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -83,7 +83,7 @@ use tokio_util::sync::CancellationToken;
|
||||||
use tracing::{debug, error, error_span, info, trace, warn};
|
use tracing::{debug, error, error_span, info, trace, warn};
|
||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
chunk_tracker::{ChunkMarkingResult, ChunkTracker},
|
chunk_tracker::{ChunkMarkingResult, ChunkTracker, HaveNeededSelected},
|
||||||
file_ops::FileOps,
|
file_ops::FileOps,
|
||||||
peer_connection::{
|
peer_connection::{
|
||||||
PeerConnection, PeerConnectionHandler, PeerConnectionOptions, WriterRequest,
|
PeerConnection, PeerConnectionHandler, PeerConnectionOptions, WriterRequest,
|
||||||
|
|
@ -203,9 +203,9 @@ impl TorrentStateLive {
|
||||||
let down_speed_estimator = SpeedEstimator::new(5);
|
let down_speed_estimator = SpeedEstimator::new(5);
|
||||||
let up_speed_estimator = SpeedEstimator::new(5);
|
let up_speed_estimator = SpeedEstimator::new(5);
|
||||||
|
|
||||||
let have_bytes = paused.have_bytes;
|
let have_bytes = paused.hns.have_bytes;
|
||||||
let needed_bytes = paused.needed_bytes;
|
let needed_bytes = paused.hns.needed_bytes;
|
||||||
let total_selected_bytes = paused.chunk_tracker.get_total_selected_bytes();
|
let total_selected_bytes = paused.hns.selected_bytes;
|
||||||
let lengths = *paused.chunk_tracker.get_lengths();
|
let lengths = *paused.chunk_tracker.get_lengths();
|
||||||
|
|
||||||
let state = Arc::new(TorrentStateLive {
|
let state = Arc::new(TorrentStateLive {
|
||||||
|
|
@ -676,8 +676,11 @@ impl TorrentStateLive {
|
||||||
files,
|
files,
|
||||||
filenames,
|
filenames,
|
||||||
chunk_tracker,
|
chunk_tracker,
|
||||||
have_bytes,
|
hns: HaveNeededSelected {
|
||||||
needed_bytes,
|
have_bytes,
|
||||||
|
needed_bytes,
|
||||||
|
selected_bytes: self.total_selected_bytes,
|
||||||
|
},
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -916,7 +919,7 @@ impl PeerHandler {
|
||||||
let n = {
|
let n = {
|
||||||
let mut n_opt = None;
|
let mut n_opt = None;
|
||||||
let bf = &live.bitfield;
|
let bf = &live.bitfield;
|
||||||
for n in g.get_chunks()?.iter_needed_pieces() {
|
for n in g.get_chunks()?.iter_queued_pieces() {
|
||||||
if bf.get(n).map(|v| *v) == Some(true) {
|
if bf.get(n).map(|v| *v) == Some(true) {
|
||||||
n_opt = Some(n);
|
n_opt = Some(n);
|
||||||
break;
|
break;
|
||||||
|
|
|
||||||
|
|
@ -353,9 +353,9 @@ impl ManagedTorrent {
|
||||||
}
|
}
|
||||||
ManagedTorrentState::Paused(p) => {
|
ManagedTorrentState::Paused(p) => {
|
||||||
resp.state = S::Paused;
|
resp.state = S::Paused;
|
||||||
resp.total_bytes = p.chunk_tracker.get_total_selected_bytes();
|
resp.total_bytes = p.hns.total();
|
||||||
resp.progress_bytes = resp.total_bytes - p.needed_bytes;
|
resp.progress_bytes = p.hns.progress();
|
||||||
resp.finished = resp.progress_bytes == resp.total_bytes;
|
resp.finished = p.hns.finished();
|
||||||
}
|
}
|
||||||
ManagedTorrentState::Live(l) => {
|
ManagedTorrentState::Live(l) => {
|
||||||
resp.state = S::Live;
|
resp.state = S::Live;
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,7 @@ use std::{collections::HashSet, fs::File, path::PathBuf, sync::Arc};
|
||||||
|
|
||||||
use parking_lot::Mutex;
|
use parking_lot::Mutex;
|
||||||
|
|
||||||
use crate::chunk_tracker::ChunkTracker;
|
use crate::chunk_tracker::{ChunkTracker, HaveNeededSelected};
|
||||||
|
|
||||||
use super::ManagedTorrentInfo;
|
use super::ManagedTorrentInfo;
|
||||||
|
|
||||||
|
|
@ -11,17 +11,15 @@ pub struct TorrentStatePaused {
|
||||||
pub(crate) files: Vec<Arc<Mutex<File>>>,
|
pub(crate) files: Vec<Arc<Mutex<File>>>,
|
||||||
pub(crate) filenames: Vec<PathBuf>,
|
pub(crate) filenames: Vec<PathBuf>,
|
||||||
pub(crate) chunk_tracker: ChunkTracker,
|
pub(crate) chunk_tracker: ChunkTracker,
|
||||||
pub(crate) have_bytes: u64,
|
pub(crate) hns: HaveNeededSelected,
|
||||||
pub(crate) needed_bytes: u64,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl TorrentStatePaused {
|
impl TorrentStatePaused {
|
||||||
pub(crate) fn update_only_files(&mut self, only_files: &HashSet<usize>) -> anyhow::Result<()> {
|
pub(crate) fn update_only_files(&mut self, only_files: &HashSet<usize>) -> anyhow::Result<()> {
|
||||||
let hn = self
|
let hns = self
|
||||||
.chunk_tracker
|
.chunk_tracker
|
||||||
.update_only_files(self.info.info.iter_file_lengths()?, only_files)?;
|
.update_only_files(self.info.info.iter_file_lengths()?, only_files)?;
|
||||||
self.have_bytes = hn.have_bytes;
|
self.hns = hns;
|
||||||
self.needed_bytes = hn.needed_bytes;
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue