All storage factories are generic now

This commit is contained in:
Igor Katson 2024-05-01 22:14:34 +01:00
parent e4adfa569a
commit 5027d8ccd1
12 changed files with 88 additions and 52 deletions

View file

@ -1,6 +1,9 @@
use std::time::Duration; use std::time::Duration;
use librqbit::{storage::mmap::MmapStorageFactory, SessionOptions}; use librqbit::{
storage::{mmap::MmapStorageFactory, StorageFactoryExt},
SessionOptions,
};
use tracing::info; use tracing::info;
#[tokio::main] #[tokio::main]
@ -28,7 +31,7 @@ async fn main() -> anyhow::Result<()> {
include_bytes!("../resources/ubuntu-21.04-live-server-amd64.iso.torrent").into(), include_bytes!("../resources/ubuntu-21.04-live-server-amd64.iso.torrent").into(),
), ),
Some(librqbit::AddTorrentOptions { Some(librqbit::AddTorrentOptions {
storage_factory: Some(Box::new(MmapStorageFactory {})), storage_factory: Some(MmapStorageFactory::default().boxed()),
paused: false, paused: false,
..Default::default() ..Default::default()
}), }),

View file

@ -16,7 +16,7 @@ use crate::{
peer_connection::PeerConnectionOptions, peer_connection::PeerConnectionOptions,
read_buf::ReadBuf, read_buf::ReadBuf,
spawn_utils::BlockingSpawner, spawn_utils::BlockingSpawner,
storage::{filesystem::FilesystemStorageFactory, StorageFactory}, storage::{filesystem::FilesystemStorageFactory, BoxStorageFactory, StorageFactoryExt},
torrent_state::{ torrent_state::{
ManagedTorrentBuilder, ManagedTorrentHandle, ManagedTorrentState, TorrentStateLive, ManagedTorrentBuilder, ManagedTorrentHandle, ManagedTorrentState, TorrentStateLive,
}, },
@ -305,7 +305,7 @@ pub struct AddTorrentOptions {
/// This is used to restore the session from serialized state. /// This is used to restore the session from serialized state.
pub preferred_id: Option<usize>, pub preferred_id: Option<usize>,
pub storage_factory: Option<Box<dyn StorageFactory>>, pub storage_factory: Option<BoxStorageFactory>,
} }
pub struct ListOnlyResponse { pub struct ListOnlyResponse {
@ -1004,7 +1004,7 @@ impl Session {
let storage_factory = opts let storage_factory = opts
.storage_factory .storage_factory
.take() .take()
.unwrap_or_else(|| Box::<FilesystemStorageFactory>::default()); .unwrap_or_else(|| FilesystemStorageFactory::default().boxed());
if opts.list_only { if opts.list_only {
return Ok(AddTorrentResponse::ListOnly(ListOnlyResponse { return Ok(AddTorrentResponse::ListOnly(ListOnlyResponse {

View file

@ -19,21 +19,21 @@ impl InMemoryPiece {
} }
} }
#[derive(Default)]
pub struct InMemoryExampleStorageFactory {} pub struct InMemoryExampleStorageFactory {}
impl StorageFactory for InMemoryExampleStorageFactory { impl StorageFactory for InMemoryExampleStorageFactory {
type Storage = InMemoryExampleStorage;
fn init_storage( fn init_storage(
&self, &self,
info: &crate::torrent_state::ManagedTorrentInfo, info: &crate::torrent_state::ManagedTorrentInfo,
) -> anyhow::Result<Box<dyn TorrentStorage>> { ) -> anyhow::Result<InMemoryExampleStorage> {
Ok(Box::new(InMemoryExampleStorage::new( InMemoryExampleStorage::new(info.lengths, info.file_infos.clone())
info.lengths,
info.file_infos.clone(),
)?))
} }
} }
struct InMemoryExampleStorage { pub struct InMemoryExampleStorage {
lengths: Lengths, lengths: Lengths,
file_infos: FileInfos, file_infos: FileInfos,
map: RwLock<HashMap<ValidPieceIndex, InMemoryPiece>>, map: RwLock<HashMap<ValidPieceIndex, InMemoryPiece>>,

View file

@ -18,7 +18,9 @@ use super::{StorageFactory, TorrentStorage};
pub struct FilesystemStorageFactory {} pub struct FilesystemStorageFactory {}
impl StorageFactory for FilesystemStorageFactory { impl StorageFactory for FilesystemStorageFactory {
fn init_storage(&self, meta: &ManagedTorrentInfo) -> anyhow::Result<Box<dyn TorrentStorage>> { type Storage = FilesystemStorage;
fn init_storage(&self, meta: &ManagedTorrentInfo) -> anyhow::Result<FilesystemStorage> {
let mut files = Vec::<OpenedFile>::new(); let mut files = Vec::<OpenedFile>::new();
let output_folder = &meta.options.output_folder; let output_folder = &meta.options.output_folder;
for file_details in meta.info.iter_file_details(&meta.lengths)? { for file_details in meta.info.iter_file_details(&meta.lengths)? {
@ -54,10 +56,10 @@ impl StorageFactory for FilesystemStorageFactory {
}; };
files.push(OpenedFile::new(file)); files.push(OpenedFile::new(file));
} }
Ok(Box::new(FilesystemStorage { Ok(FilesystemStorage {
output_folder: output_folder.clone(), output_folder: output_folder.clone(),
opened_files: files, opened_files: files,
})) })
} }
} }

View file

@ -6,26 +6,26 @@ use crate::{FileInfos, ManagedTorrentInfo};
use super::{StorageFactory, TorrentStorage}; use super::{StorageFactory, TorrentStorage};
#[derive(Default)]
pub struct MmapStorageFactory {} pub struct MmapStorageFactory {}
struct MmapStorage { pub struct MmapStorage {
mmap: RwLock<MmapMut>, mmap: RwLock<MmapMut>,
file_infos: FileInfos, file_infos: FileInfos,
} }
impl StorageFactory for MmapStorageFactory { impl StorageFactory for MmapStorageFactory {
fn init_storage( type Storage = MmapStorage;
&self,
info: &ManagedTorrentInfo, fn init_storage(&self, info: &ManagedTorrentInfo) -> anyhow::Result<Self::Storage> {
) -> anyhow::Result<Box<dyn crate::storage::TorrentStorage>> { Ok(MmapStorage {
Ok(Box::new(MmapStorage {
mmap: RwLock::new( mmap: RwLock::new(
MmapOptions::new() MmapOptions::new()
.len(info.lengths.total_length().try_into()?) .len(info.lengths.total_length().try_into()?)
.map_anon()?, .map_anon()?,
), ),
file_infos: info.file_infos.clone(), file_infos: info.file_infos.clone(),
})) })
} }
} }

View file

@ -9,11 +9,40 @@ use std::{any::Any, path::Path};
use crate::torrent_state::ManagedTorrentInfo; use crate::torrent_state::ManagedTorrentInfo;
pub trait StorageFactory: Send + Sync + Any { pub trait StorageFactory: Send + Sync + Any {
fn init_storage(&self, info: &ManagedTorrentInfo) -> anyhow::Result<Box<dyn TorrentStorage>>; type Storage: TorrentStorage;
fn init_storage(&self, info: &ManagedTorrentInfo) -> anyhow::Result<Self::Storage>;
}
pub type BoxStorageFactory = Box<dyn StorageFactory<Storage = Box<dyn TorrentStorage>>>;
pub trait StorageFactoryExt {
fn boxed(self) -> BoxStorageFactory;
}
impl<SF: StorageFactory> StorageFactoryExt for SF {
fn boxed(self) -> BoxStorageFactory {
struct BoxFactory<SF> {
sf: SF,
}
impl<SF: StorageFactory> StorageFactory for BoxFactory<SF> {
type Storage = Box<dyn TorrentStorage>;
fn init_storage(&self, info: &ManagedTorrentInfo) -> anyhow::Result<Self::Storage> {
let s = self.sf.init_storage(info)?;
Ok(Box::new(s))
}
}
Box::new(BoxFactory { sf: self })
}
} }
impl<U: StorageFactory + ?Sized> StorageFactory for Box<U> { impl<U: StorageFactory + ?Sized> StorageFactory for Box<U> {
fn init_storage(&self, info: &ManagedTorrentInfo) -> anyhow::Result<Box<dyn TorrentStorage>> { type Storage = U::Storage;
fn init_storage(&self, info: &ManagedTorrentInfo) -> anyhow::Result<U::Storage> {
(**self).init_storage(info) (**self).init_storage(info)
} }
} }

View file

@ -17,17 +17,16 @@ impl<U: StorageFactory> SlowStorageFactory<U> {
} }
impl<U: StorageFactory> StorageFactory for SlowStorageFactory<U> { impl<U: StorageFactory> StorageFactory for SlowStorageFactory<U> {
fn init_storage( type Storage = SlowStorage<U::Storage>;
&self,
info: &crate::ManagedTorrentInfo, fn init_storage(&self, info: &crate::ManagedTorrentInfo) -> anyhow::Result<Self::Storage> {
) -> anyhow::Result<Box<dyn TorrentStorage>> { Ok(SlowStorage {
Ok(Box::new(SlowStorage {
underlying: self.underlying_factory.init_storage(info)?, underlying: self.underlying_factory.init_storage(info)?,
})) })
} }
} }
struct SlowStorage<U> { pub struct SlowStorage<U> {
underlying: U, underlying: U,
} }

View file

@ -15,18 +15,17 @@ impl<U> TimingStorageFactory<U> {
} }
impl<U: StorageFactory> StorageFactory for TimingStorageFactory<U> { impl<U: StorageFactory> StorageFactory for TimingStorageFactory<U> {
fn init_storage( type Storage = TimingStorage<U::Storage>;
&self,
info: &crate::ManagedTorrentInfo, fn init_storage(&self, info: &crate::ManagedTorrentInfo) -> anyhow::Result<Self::Storage> {
) -> anyhow::Result<Box<dyn TorrentStorage>> { Ok(TimingStorage {
Ok(Box::new(TimingStorage {
name: self.name.clone(), name: self.name.clone(),
underlying: self.underlying_factory.init_storage(info)?, underlying: self.underlying_factory.init_storage(info)?,
})) })
} }
} }
struct TimingStorage<U> { pub struct TimingStorage<U> {
name: String, name: String,
underlying: U, underlying: U,
} }

View file

@ -5,8 +5,9 @@ use tokio::{io::AsyncReadExt, time::timeout};
use tracing::info; use tracing::info;
use crate::{ use crate::{
create_torrent, storage::example::InMemoryExampleStorageFactory, AddTorrent, create_torrent,
CreateTorrentOptions, Session, storage::{example::InMemoryExampleStorageFactory, StorageFactoryExt},
AddTorrent, CreateTorrentOptions, Session,
}; };
use super::test_util::create_default_random_dir_with_torrents; use super::test_util::create_default_random_dir_with_torrents;
@ -85,7 +86,7 @@ async fn e2e_stream() -> anyhow::Result<()> {
AddTorrent::from_bytes(torrent.as_bytes()?), AddTorrent::from_bytes(torrent.as_bytes()?),
Some(crate::AddTorrentOptions { Some(crate::AddTorrentOptions {
paused: false, paused: false,
storage_factory: Some(Box::new(InMemoryExampleStorageFactory {})), storage_factory: Some(InMemoryExampleStorageFactory::default().boxed()),
initial_peers: Some(vec![peer]), initial_peers: Some(vec![peer]),
..Default::default() ..Default::default()
}), }),

View file

@ -8,7 +8,11 @@ use anyhow::Context;
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, storage::StorageFactory}; use crate::{
chunk_tracker::ChunkTracker,
file_ops::FileOps,
storage::{BoxStorageFactory, StorageFactory},
};
use super::{paused::TorrentStatePaused, ManagedTorrentInfo}; use super::{paused::TorrentStatePaused, ManagedTorrentInfo};
@ -34,7 +38,7 @@ impl TorrentStateInitializing {
pub async fn check( pub async fn check(
&self, &self,
storage_factory: &dyn StorageFactory, storage_factory: &BoxStorageFactory,
) -> anyhow::Result<TorrentStatePaused> { ) -> anyhow::Result<TorrentStatePaused> {
let files = storage_factory.init_storage(&self.meta)?; let files = storage_factory.init_storage(&self.meta)?;
info!("Doing initial checksum validation, this might take a while..."); info!("Doing initial checksum validation, this might take a while...");

View file

@ -36,7 +36,7 @@ use tracing::warn;
use crate::chunk_tracker::ChunkTracker; use crate::chunk_tracker::ChunkTracker;
use crate::file_info::FileInfo; use crate::file_info::FileInfo;
use crate::spawn_utils::BlockingSpawner; use crate::spawn_utils::BlockingSpawner;
use crate::storage::StorageFactory; use crate::storage::BoxStorageFactory;
use crate::torrent_state::stats::LiveStats; use crate::torrent_state::stats::LiveStats;
use crate::type_aliases::FileInfos; use crate::type_aliases::FileInfos;
use crate::type_aliases::PeerStream; use crate::type_aliases::PeerStream;
@ -108,7 +108,7 @@ pub struct ManagedTorrentInfo {
pub struct ManagedTorrent { pub struct ManagedTorrent {
pub info: Arc<ManagedTorrentInfo>, pub info: Arc<ManagedTorrentInfo>,
pub(crate) storage_factory: Box<dyn StorageFactory>, pub(crate) storage_factory: BoxStorageFactory,
state_change_notify: Notify, state_change_notify: Notify,
locked: RwLock<ManagedTorrentLocked>, locked: RwLock<ManagedTorrentLocked>,
@ -273,7 +273,7 @@ impl ManagedTorrent {
error_span!(parent: span.clone(), "initialize_and_start"), error_span!(parent: span.clone(), "initialize_and_start"),
token.clone(), token.clone(),
async move { async move {
match init.check(&*t.storage_factory).await { match init.check(&t.storage_factory).await {
Ok(paused) => { Ok(paused) => {
let mut g = t.locked.write(); let mut g = t.locked.write();
if let ManagedTorrentState::Initializing(_) = &g.state { if let ManagedTorrentState::Initializing(_) = &g.state {
@ -504,7 +504,7 @@ pub(crate) struct ManagedTorrentBuilder {
peer_id: Option<Id20>, peer_id: Option<Id20>,
spawner: Option<BlockingSpawner>, spawner: Option<BlockingSpawner>,
allow_overwrite: bool, allow_overwrite: bool,
storage_factory: Box<dyn StorageFactory>, storage_factory: BoxStorageFactory,
} }
impl ManagedTorrentBuilder { impl ManagedTorrentBuilder {
@ -512,7 +512,7 @@ impl ManagedTorrentBuilder {
info: TorrentMetaV1Info<ByteBufOwned>, info: TorrentMetaV1Info<ByteBufOwned>,
info_hash: Id20, info_hash: Id20,
output_folder: PathBuf, output_folder: PathBuf,
storage_factory: Box<dyn StorageFactory>, storage_factory: BoxStorageFactory,
) -> Self { ) -> Self {
Self { Self {
info, info,

View file

@ -8,8 +8,7 @@ use librqbit::{
http_api::{HttpApi, HttpApiOptions}, http_api::{HttpApi, HttpApiOptions},
http_api_client, librqbit_spawn, http_api_client, librqbit_spawn,
storage::{ storage::{
filesystem::FilesystemStorageFactory, slow::SlowStorageFactory, filesystem::FilesystemStorageFactory, timing::TimingStorageFactory, StorageFactoryExt,
timing::TimingStorageFactory,
}, },
tracing_subscriber_config_utils::{init_logging, InitLoggingOptions}, tracing_subscriber_config_utils::{init_logging, InitLoggingOptions},
AddTorrent, AddTorrentOptions, AddTorrentResponse, Api, ListOnlyResponse, AddTorrent, AddTorrentOptions, AddTorrentResponse, Api, ListOnlyResponse,
@ -382,9 +381,9 @@ async fn async_main(opts: Opts) -> anyhow::Result<()> {
initial_peers: download_opts.initial_peers.clone().map(|p| p.0), initial_peers: download_opts.initial_peers.clone().map(|p| p.0),
disable_trackers: download_opts.disable_trackers, disable_trackers: download_opts.disable_trackers,
storage_factory: Some({ storage_factory: Some({
let sf = Box::<FilesystemStorageFactory>::default(); let sf = FilesystemStorageFactory::default();
// let sf = Box::new(SlowStorageFactory::new(sf)); // let sf = SlowStorageFactory::new(sf);
Box::new(TimingStorageFactory::new("fs".to_owned(), sf)) TimingStorageFactory::new("fs".to_owned(), sf).boxed()
}), }),
..Default::default() ..Default::default()
}; };