This commit is contained in:
Igor Katson 2021-07-01 19:17:44 +01:00
parent a6981231c1
commit 5942e6a9d5
12 changed files with 186 additions and 148 deletions

44
Cargo.lock generated
View file

@ -148,24 +148,6 @@ dependencies = [
"syn", "syn",
] ]
[[package]]
name = "commoncrypto"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d056a8586ba25a1e4d61cb090900e495952c7886786fc55f909ab2f819b69007"
dependencies = [
"commoncrypto-sys",
]
[[package]]
name = "commoncrypto-sys"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1fed34f46747aa73dfaa578069fd8279d2818ade2b55f38f22a9401c7f4083e2"
dependencies = [
"libc",
]
[[package]] [[package]]
name = "core-foundation" name = "core-foundation"
version = "0.9.1" version = "0.9.1"
@ -191,18 +173,6 @@ dependencies = [
"libc", "libc",
] ]
[[package]]
name = "crypto-hash"
version = "0.3.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8a77162240fd97248d19a564a565eb563a3f592b386e4136fb300909e67dddca"
dependencies = [
"commoncrypto",
"hex",
"openssl",
"winapi",
]
[[package]] [[package]]
name = "digest" name = "digest"
version = "0.9.0" version = "0.9.0"
@ -474,12 +444,6 @@ dependencies = [
"libc", "libc",
] ]
[[package]]
name = "hex"
version = "0.3.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "805026a5d0141ffc30abb3be3173848ad46a1b1664fe632428479619a3644d77"
[[package]] [[package]]
name = "http" name = "http"
version = "0.2.4" version = "0.2.4"
@ -640,7 +604,6 @@ dependencies = [
"bincode", "bincode",
"bitvec", "bitvec",
"byteorder", "byteorder",
"crypto-hash",
"futures", "futures",
"log", "log",
"openssl", "openssl",
@ -648,7 +611,6 @@ dependencies = [
"rand 0.8.4", "rand 0.8.4",
"reqwest", "reqwest",
"serde", "serde",
"sha1",
"size_format", "size_format",
"tokio", "tokio",
"urlencoding", "urlencoding",
@ -1323,12 +1285,6 @@ dependencies = [
"opaque-debug", "opaque-debug",
] ]
[[package]]
name = "sha1"
version = "0.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2579985fda508104f7587689507983eadd6a6e84dd35d6d115361f530916fa0d"
[[package]] [[package]]
name = "size_format" name = "size_format"
version = "1.0.2" version = "1.0.2"

View file

@ -4,7 +4,8 @@ version = "0.1.0"
authors = ["Igor Katson <igor.katson@gmail.com>"] authors = ["Igor Katson <igor.katson@gmail.com>"]
edition = "2018" edition = "2018"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html [features]
rt-single-thread = []
[dependencies] [dependencies]
librqbit = {path="./crates/librqbit"} librqbit = {path="./crates/librqbit"}

View file

@ -3,6 +3,13 @@ name = "librqbit"
version = "0.1.0" version = "0.1.0"
authors = ["Igor Katson <igor.katson@gmail.com>"] authors = ["Igor Katson <igor.katson@gmail.com>"]
edition = "2018" edition = "2018"
default-features = ["sha1-system"]
[features]
default = ["sha1-openssl"]
sha1-system = ["crypto-hash"]
sha1-openssl = ["openssl"]
sha1-rust = ["sha1"]
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
@ -10,7 +17,7 @@ edition = "2018"
tokio = {version = "1", features = ["macros", "rt-multi-thread"]} tokio = {version = "1", features = ["macros", "rt-multi-thread"]}
serde = {version = "1", features=["derive"]} serde = {version = "1", features=["derive"]}
anyhow = "1" anyhow = "1"
sha1 = "0.6"
reqwest = "0.11" reqwest = "0.11"
urlencoding = "1" urlencoding = "1"
byteorder = "1" byteorder = "1"
@ -20,9 +27,11 @@ parking_lot = "0.11"
log = "0.4" log = "0.4"
size_format = "1" size_format = "1"
rand = "0.8" rand = "0.8"
openssl = "0.10"
warp = "0.3" warp = "0.3"
crypto-hash = "0.3"
openssl = {version="0.10", optional=true}
crypto-hash = {version="0.3", optional=true}
sha1 = {version = "0.6", optional=true}
uuid = {version = "0.8", features = ["v4"]} uuid = {version = "0.8", features = ["v4"]}
futures = "0.3" futures = "0.3"

View file

@ -1,6 +1,7 @@
use std::{ use std::{
fs::File, fs::File,
io::{Read, Seek, SeekFrom, Write}, io::{Read, Seek, SeekFrom, Write},
marker::PhantomData,
sync::Arc, sync::Arc,
}; };
@ -12,7 +13,7 @@ use crate::{
buffers::ByteString, buffers::ByteString,
lengths::{ChunkInfo, Lengths, ValidPieceIndex}, lengths::{ChunkInfo, Lengths, ValidPieceIndex},
peer_binary_protocol::Piece, peer_binary_protocol::Piece,
sha1w::{self, ISha1}, sha1w::ISha1,
torrent_metainfo::{FileIteratorName, TorrentMetaV1Owned}, torrent_metainfo::{FileIteratorName, TorrentMetaV1Owned},
type_aliases::{PeerHandle, BF}, type_aliases::{PeerHandle, BF},
}; };
@ -24,6 +25,7 @@ pub struct InitialCheckResults {
pub needed_bytes: u64, pub needed_bytes: u64,
} }
#[inline(never)]
pub fn update_hash_from_file<Sha1: ISha1>( pub fn update_hash_from_file<Sha1: ISha1>(
file: &mut File, file: &mut File,
hash: &mut Sha1, hash: &mut Sha1,
@ -46,13 +48,14 @@ pub fn update_hash_from_file<Sha1: ISha1>(
Ok(()) Ok(())
} }
pub struct FileOps<'a> { pub struct FileOps<'a, Sha1> {
torrent: &'a TorrentMetaV1Owned, torrent: &'a TorrentMetaV1Owned,
files: &'a [Arc<Mutex<File>>], files: &'a [Arc<Mutex<File>>],
lengths: &'a Lengths, lengths: &'a Lengths,
phantom_data: PhantomData<Sha1>,
} }
impl<'a> FileOps<'a> { impl<'a, Sha1Impl: ISha1> FileOps<'a, Sha1Impl> {
pub fn new( pub fn new(
torrent: &'a TorrentMetaV1Owned, torrent: &'a TorrentMetaV1Owned,
files: &'a [Arc<Mutex<File>>], files: &'a [Arc<Mutex<File>>],
@ -62,6 +65,7 @@ impl<'a> FileOps<'a> {
torrent, torrent,
files, files,
lengths, lengths,
phantom_data: PhantomData,
} }
} }
@ -121,7 +125,7 @@ impl<'a> FileOps<'a> {
let mut read_buffer = vec![0u8; 65536]; let mut read_buffer = vec![0u8; 65536];
for piece_info in self.lengths.iter_piece_infos() { for piece_info in self.lengths.iter_piece_infos() {
let mut computed_hash = sha1w::Sha1System::new(); let mut computed_hash = Sha1Impl::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 at_least_one_file_required = current_file.full_file_required;
@ -215,13 +219,14 @@ impl<'a> FileOps<'a> {
}) })
} }
#[inline(never)]
pub fn check_piece( pub fn check_piece(
&self, &self,
who_sent: PeerHandle, who_sent: PeerHandle,
piece_index: ValidPieceIndex, piece_index: ValidPieceIndex,
last_received_chunk: &ChunkInfo, last_received_chunk: &ChunkInfo,
) -> anyhow::Result<bool> { ) -> anyhow::Result<bool> {
let mut h = sha1w::Sha1System::new(); let mut h = Sha1Impl::new();
let piece_length = self.lengths.piece_length(piece_index); let piece_length = self.lengths.piece_length(piece_index);
let mut absolute_offset = self.lengths.piece_offset(piece_index); let mut absolute_offset = self.lengths.piece_offset(piece_index);
let mut buf = vec![0u8; std::cmp::min(65536, piece_length as usize)]; let mut buf = vec![0u8; std::cmp::min(65536, piece_length as usize)];

View file

@ -19,7 +19,7 @@ use crate::{
}, },
peer_id::try_decode_peer_id, peer_id::try_decode_peer_id,
peer_state::InflightRequest, peer_state::InflightRequest,
spawn_utils::{spawn, spawn_block_in_place}, spawn_utils::{spawn, BlockingSpawner},
torrent_state::TorrentState, torrent_state::TorrentState,
type_aliases::PeerHandle, type_aliases::PeerHandle,
}; };
@ -33,11 +33,12 @@ pub enum WriterRequest {
#[derive(Clone)] #[derive(Clone)]
pub struct PeerConnection { pub struct PeerConnection {
state: Arc<TorrentState>, state: Arc<TorrentState>,
spawner: BlockingSpawner,
} }
impl PeerConnection { impl PeerConnection {
pub fn new(state: Arc<TorrentState>) -> Self { pub fn new(state: Arc<TorrentState>, spawner: BlockingSpawner) -> Self {
PeerConnection { state } PeerConnection { state, spawner }
} }
pub fn into_state(self) -> Arc<TorrentState> { pub fn into_state(self) -> Arc<TorrentState> {
self.state self.state
@ -128,14 +129,15 @@ impl PeerConnection {
let preamble_len = serialize_piece_preamble(&chunk, &mut buf); let preamble_len = serialize_piece_preamble(&chunk, &mut buf);
let full_len = preamble_len + chunk.size as usize; let full_len = preamble_len + chunk.size as usize;
buf.resize(full_len, 0); buf.resize(full_len, 0);
spawn_block_in_place(|| { this.spawner
this.state.file_ops().read_chunk( .spawn_block_in_place(|| {
handle, this.state.file_ops().read_chunk(
&chunk, handle,
&mut buf[preamble_len..], &chunk,
) &mut buf[preamble_len..],
}) )
.with_context(|| format!("error reading chunk {:?}", chunk))?; })
.with_context(|| format!("error reading chunk {:?}", chunk))?;
uploaded_add = Some(chunk.size); uploaded_add = Some(chunk.size);
full_len full_len
} }
@ -544,79 +546,81 @@ impl PeerConnection {
// to prevent deadlocks. // to prevent deadlocks.
drop(g); drop(g);
spawn_block_in_place(move || { self.spawner
let index = piece.index; .spawn_block_in_place(move || {
let index = piece.index;
// TODO: in theory we should unmark the piece as downloaded here. But if there was a disk error, what // TODO: in theory we should unmark the piece as downloaded here. But if there was a disk error, what
// should we really do? If we unmark it, it will get requested forever... // should we really do? If we unmark it, it will get requested forever...
// //
// So let's just unwrap and abort. // So let's just unwrap and abort.
self.state self.state
.file_ops() .file_ops()
.write_chunk(handle, &piece, &chunk_info) .write_chunk(handle, &piece, &chunk_info)
.expect("expected to be able to write to disk"); .expect("expected to be able to write to disk");
let full_piece_download_time = match full_piece_download_time { let full_piece_download_time = match full_piece_download_time {
Some(t) => t, Some(t) => t,
None => return Ok(()), None => return Ok(()),
}; };
match self match self
.state .state
.file_ops() .file_ops()
.check_piece(handle, chunk_info.piece_index, &chunk_info) .check_piece(handle, chunk_info.piece_index, &chunk_info)
.with_context(|| format!("error checking piece={}", index))? .with_context(|| format!("error checking piece={}", index))?
{ {
true => { true => {
let piece_len = self.state.lengths.piece_length(chunk_info.piece_index) as u64; let piece_len =
self.state self.state.lengths.piece_length(chunk_info.piece_index) as u64;
.stats self.state
.downloaded_and_checked .stats
.fetch_add(piece_len, Ordering::Relaxed); .downloaded_and_checked
self.state .fetch_add(piece_len, Ordering::Relaxed);
.stats self.state
.have .stats
.fetch_add(piece_len, Ordering::Relaxed); .have
self.state .fetch_add(piece_len, Ordering::Relaxed);
.stats self.state
.downloaded_pieces .stats
.fetch_add(1, Ordering::Relaxed); .downloaded_pieces
self.state .fetch_add(1, Ordering::Relaxed);
.stats self.state
.downloaded_pieces .stats
.fetch_add(1, Ordering::Relaxed); .downloaded_pieces
self.state.stats.total_piece_download_ms.fetch_add( .fetch_add(1, Ordering::Relaxed);
full_piece_download_time.as_millis() as u64, self.state.stats.total_piece_download_ms.fetch_add(
Ordering::Relaxed, full_piece_download_time.as_millis() as u64,
); Ordering::Relaxed,
self.state );
.locked self.state
.write() .locked
.chunks .write()
.mark_piece_downloaded(chunk_info.piece_index); .chunks
.mark_piece_downloaded(chunk_info.piece_index);
debug!( debug!(
"piece={} successfully downloaded and verified from {}", "piece={} successfully downloaded and verified from {}",
index, handle index, handle
); );
self.state.maybe_transmit_haves(chunk_info.piece_index); self.state.maybe_transmit_haves(chunk_info.piece_index);
} }
false => { false => {
warn!( warn!(
"checksum for piece={} did not validate, came from {}", "checksum for piece={} did not validate, came from {}",
index, handle index, handle
); );
self.state self.state
.locked .locked
.write() .write()
.chunks .chunks
.mark_piece_broken(chunk_info.piece_index); .mark_piece_broken(chunk_info.piece_index);
} }
}; };
Ok::<_, anyhow::Error>(()) Ok::<_, anyhow::Error>(())
}) })
.with_context(|| format!("error processing received chunk {:?}", chunk_info))?; .with_context(|| format!("error processing received chunk {:?}", chunk_info))?;
Ok(()) Ok(())
} }
} }

View file

@ -7,6 +7,8 @@ use std::marker::PhantomData;
use crate::buffers::ByteBuf; use crate::buffers::ByteBuf;
use crate::buffers::ByteString; use crate::buffers::ByteString;
use crate::clone_to_owned::CloneToOwned; use crate::clone_to_owned::CloneToOwned;
use crate::sha1w::ISha1;
use crate::type_aliases::Sha1;
pub struct BencodeDeserializer<'de> { pub struct BencodeDeserializer<'de> {
buf: &'de [u8], buf: &'de [u8],
@ -536,9 +538,9 @@ impl<'a, 'de> serde::de::MapAccess<'de> for MapAccess<'a, 'de> {
let value = seed.deserialize(&mut *self.de)?; let value = seed.deserialize(&mut *self.de)?;
if self.de.is_torrent_info && self.de.field_context.as_slice() == [ByteBuf(b"info")] { if self.de.is_torrent_info && self.de.field_context.as_slice() == [ByteBuf(b"info")] {
let len = self.de.buf.as_ptr() as usize - buf_before.as_ptr() as usize; let len = self.de.buf.as_ptr() as usize - buf_before.as_ptr() as usize;
let mut hash = sha1::Sha1::new(); let mut hash = Sha1::new();
hash.update(&buf_before[..len]); hash.update(&buf_before[..len]);
let digest = hash.digest().bytes(); let digest = hash.finish();
self.de.torrent_info_digest = Some(digest) self.de.torrent_info_digest = Some(digest)
} }
self.de.field_context.pop(); self.de.field_context.pop();

View file

@ -12,10 +12,12 @@ pub trait ISha1 {
fn finish(self) -> [u8; 20]; fn finish(self) -> [u8; 20];
} }
#[cfg(feature = "sha1-rust")]
pub struct Sha1Rust { pub struct Sha1Rust {
inner: sha1::Sha1, inner: sha1::Sha1,
} }
#[cfg(feature = "sha1-rust")]
impl ISha1 for Sha1Rust { impl ISha1 for Sha1Rust {
fn new() -> Self { fn new() -> Self {
Sha1Rust { Sha1Rust {
@ -31,9 +33,13 @@ impl ISha1 for Sha1Rust {
self.inner.digest().bytes() self.inner.digest().bytes()
} }
} }
#[cfg(feature = "sha1-openssl")]
pub struct Sha1Openssl { pub struct Sha1Openssl {
inner: openssl::sha::Sha1, inner: openssl::sha::Sha1,
} }
#[cfg(feature = "sha1-openssl")]
impl ISha1 for Sha1Openssl { impl ISha1 for Sha1Openssl {
fn new() -> Self { fn new() -> Self {
Self { Self {
@ -50,10 +56,12 @@ impl ISha1 for Sha1Openssl {
} }
} }
#[cfg(feature = "sha1-system")]
pub struct Sha1System { pub struct Sha1System {
inner: crypto_hash::Hasher, inner: crypto_hash::Hasher,
} }
#[cfg(feature = "sha1-system")]
impl ISha1 for Sha1System { impl ISha1 for Sha1System {
fn new() -> Self { fn new() -> Self {
Self { Self {

View file

@ -19,9 +19,22 @@ pub fn spawn<N: Display + 'static + Send>(
}); });
} }
pub fn spawn_block_in_place<F: FnOnce() -> R, R>(f: F) -> R { #[derive(Clone, Copy, Debug)]
// Have this wrapper so that it's easy to switch to just f() when pub struct BlockingSpawner {
// using tokio's single-threaded runtime. Single-threaded runtime is allow_tokio_block_in_place: bool,
// easier to read with time profilers. }
tokio::task::block_in_place(f)
impl BlockingSpawner {
pub fn new(allow_tokio_block_in_place: bool) -> Self {
Self {
allow_tokio_block_in_place,
}
}
pub fn spawn_block_in_place<F: FnOnce() -> R, R>(&self, f: F) -> R {
if self.allow_tokio_block_in_place {
return tokio::task::block_in_place(f);
}
return f();
}
} }

View file

@ -21,11 +21,12 @@ use crate::{
file_ops::FileOps, file_ops::FileOps,
http_api::make_and_run_http_api, http_api::make_and_run_http_api,
lengths::Lengths, lengths::Lengths,
spawn_utils::spawn, spawn_utils::{spawn, BlockingSpawner},
speed_estimator::SpeedEstimator, speed_estimator::SpeedEstimator,
torrent_metainfo::TorrentMetaV1Owned, torrent_metainfo::TorrentMetaV1Owned,
torrent_state::{AtomicStats, TorrentState, TorrentStateLocked}, torrent_state::{AtomicStats, TorrentState, TorrentStateLocked},
tracker_comms::{TrackerError, TrackerRequest, TrackerRequestEvent, TrackerResponse}, tracker_comms::{TrackerError, TrackerRequest, TrackerRequestEvent, TrackerResponse},
type_aliases::Sha1,
}; };
pub struct TorrentManagerBuilder { pub struct TorrentManagerBuilder {
torrent: TorrentMetaV1Owned, torrent: TorrentMetaV1Owned,
@ -33,6 +34,7 @@ pub struct TorrentManagerBuilder {
output_folder: PathBuf, output_folder: PathBuf,
only_files: Option<Vec<usize>>, only_files: Option<Vec<usize>>,
force_tracker_interval: Option<Duration>, force_tracker_interval: Option<Duration>,
spawner: Option<BlockingSpawner>,
} }
impl TorrentManagerBuilder { impl TorrentManagerBuilder {
@ -43,6 +45,7 @@ impl TorrentManagerBuilder {
output_folder: output_folder.as_ref().into(), output_folder: output_folder.as_ref().into(),
only_files: None, only_files: None,
force_tracker_interval: None, force_tracker_interval: None,
spawner: None,
} }
} }
@ -61,6 +64,11 @@ impl TorrentManagerBuilder {
self self
} }
pub fn spawner(&mut self, spawner: BlockingSpawner) -> &mut Self {
self.spawner = Some(spawner);
self
}
pub async fn start_manager(self) -> anyhow::Result<TorrentManagerHandle> { pub async fn start_manager(self) -> anyhow::Result<TorrentManagerHandle> {
TorrentManager::start( TorrentManager::start(
self.torrent, self.torrent,
@ -68,6 +76,7 @@ impl TorrentManagerBuilder {
self.overwrite, self.overwrite,
self.only_files, self.only_files,
self.force_tracker_interval, self.force_tracker_interval,
self.spawner.unwrap_or(BlockingSpawner::new(true)),
) )
} }
} }
@ -114,6 +123,7 @@ impl TorrentManager {
overwrite: bool, overwrite: bool,
only_files: Option<Vec<usize>>, only_files: Option<Vec<usize>>,
force_tracker_interval: Option<Duration>, force_tracker_interval: Option<Duration>,
spawner: BlockingSpawner,
) -> anyhow::Result<TorrentManagerHandle> { ) -> anyhow::Result<TorrentManagerHandle> {
let files = { let files = {
let mut files = let mut files =
@ -155,8 +165,8 @@ impl TorrentManager {
debug!("computed lengths: {:?}", &lengths); debug!("computed lengths: {:?}", &lengths);
info!("Doing initial checksum validation, this might take a while..."); info!("Doing initial checksum validation, this might take a while...");
let initial_check_results = let initial_check_results = FileOps::<Sha1>::new(&torrent, &files, &lengths)
FileOps::new(&torrent, &files, &lengths).initial_check(only_files.as_deref())?; .initial_check(only_files.as_deref())?;
info!( info!(
"Initial check results: have {}, needed {}", "Initial check results: have {}, needed {}",
@ -185,6 +195,7 @@ impl TorrentManager {
}, },
needed: initial_check_results.needed_bytes, needed: initial_check_results.needed_bytes,
lengths, lengths,
spawner,
}); });
let estimator = Arc::new(SpeedEstimator::new(5)); let estimator = Arc::new(SpeedEstimator::new(5));

View file

@ -21,9 +21,9 @@ use crate::{
peer_binary_protocol::{Handshake, Message}, peer_binary_protocol::{Handshake, Message},
peer_connection::{PeerConnection, WriterRequest}, peer_connection::{PeerConnection, WriterRequest},
peer_state::{LivePeerState, PeerState}, peer_state::{LivePeerState, PeerState},
spawn_utils::spawn, spawn_utils::{spawn, BlockingSpawner},
torrent_metainfo::TorrentMetaV1Owned, torrent_metainfo::TorrentMetaV1Owned,
type_aliases::{PeerHandle, BF}, type_aliases::{PeerHandle, Sha1, BF},
}; };
pub struct InflightPiece { pub struct InflightPiece {
@ -192,10 +192,12 @@ pub struct TorrentState {
pub lengths: Lengths, pub lengths: Lengths,
pub needed: u64, pub needed: u64,
pub stats: AtomicStats, pub stats: AtomicStats,
pub spawner: BlockingSpawner,
} }
impl TorrentState { impl TorrentState {
pub fn file_ops(&self) -> FileOps<'_> { pub fn file_ops(&self) -> FileOps<'_, Sha1> {
FileOps::new(&self.torrent, &self.files, &self.lengths) FileOps::new(&self.torrent, &self.files, &self.lengths)
} }
@ -400,7 +402,7 @@ impl TorrentState {
None => return false, None => return false,
}; };
let peer_connection = PeerConnection::new(self.clone()); let peer_connection = PeerConnection::new(self.clone(), self.spawner.clone());
spawn(format!("manage_peer({})", handle), async move { spawn(format!("manage_peer({})", handle), async move {
if let Err(e) = peer_connection.manage_peer(addr, handle, out_rx).await { if let Err(e) = peer_connection.manage_peer(addr, handle, out_rx).await {
debug!("error managing peer {}: {:#}", handle, e) debug!("error managing peer {}: {:#}", handle, e)

View file

@ -3,3 +3,12 @@ use std::net::SocketAddr;
pub type BF = bitvec::vec::BitVec<bitvec::order::Msb0, u8>; pub type BF = bitvec::vec::BitVec<bitvec::order::Msb0, u8>;
pub type PeerHandle = SocketAddr; pub type PeerHandle = SocketAddr;
#[cfg(feature = "sha1-openssl")]
pub type Sha1 = crate::sha1w::Sha1Openssl;
#[cfg(feature = "sha1-system")]
pub type Sha1 = crate::sha1w::Sha1System;
#[cfg(feature = "sha1-rust")]
pub type Sha1 = crate::sha1w::Sha1Rust;

View file

@ -3,6 +3,7 @@ use std::{fs::File, io::Read, time::Duration};
use anyhow::Context; use anyhow::Context;
use clap::Clap; use clap::Clap;
use librqbit::{ use librqbit::{
spawn_utils::BlockingSpawner,
torrent_manager::TorrentManagerBuilder, torrent_manager::TorrentManagerBuilder,
torrent_metainfo::{torrent_from_bytes, TorrentMetaV1Owned}, torrent_metainfo::{torrent_from_bytes, TorrentMetaV1Owned},
}; };
@ -77,6 +78,12 @@ struct Opts {
/// pretty big, e.g. 30 minutes. This can force a certain value. /// pretty big, e.g. 30 minutes. This can force a certain value.
#[clap(short = 'i', long = "tracker-refresh-interval")] #[clap(short = 'i', long = "tracker-refresh-interval")]
force_tracker_interval: Option<u64>, force_tracker_interval: Option<u64>,
/// Set this flag if you want to use tokio's single threaded runtime.
/// It MAY perform better, but the main purpose is easier debugging, as time
/// profilers work better with this one.
#[clap(short, long)]
single_thread_runtime: bool,
} }
fn compute_only_files( fn compute_only_files(
@ -129,7 +136,18 @@ fn main() -> anyhow::Result<()> {
init_logging(&opts); init_logging(&opts);
let rt = tokio::runtime::Builder::new_multi_thread() let (mut rt_builder, spawner) = match opts.single_thread_runtime {
true => (
tokio::runtime::Builder::new_current_thread(),
BlockingSpawner::new(false),
),
false => (
tokio::runtime::Builder::new_multi_thread(),
BlockingSpawner::new(true),
),
};
let rt = rt_builder
.enable_time() .enable_time()
.enable_io() .enable_io()
// the default is 512, it can get out of hand, as this program is CPU-bound on // the default is 512, it can get out of hand, as this program is CPU-bound on
@ -161,7 +179,7 @@ fn main() -> anyhow::Result<()> {
}; };
let mut builder = TorrentManagerBuilder::new(torrent, opts.output_folder); let mut builder = TorrentManagerBuilder::new(torrent, opts.output_folder);
builder.overwrite(opts.overwrite); builder.overwrite(opts.overwrite).spawner(spawner);
if let Some(only_files) = only_files { if let Some(only_files) = only_files {
builder.only_files(only_files); builder.only_files(only_files);
} }