Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 57 additions & 15 deletions crates/backend/src/opendal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,14 @@ use crate::reqwest::reqwest_client;
mod constants {
/// Default number of retries
pub(super) const DEFAULT_RETRY: usize = 5;

/// B2 `b2_list_file_names` page size to request.
///
/// `OpenDAL` only sends `maxFileCount` when `ListOptions.limit` is set. If we
/// omit it, B2 defaults to 100 names per page (max 10000). Prune (and
/// check) list every pack under `data/`, so the default turns a few dozen
/// round-trips into thousands and dominates runtime against B2.
pub(super) const B2_LIST_PAGE_SIZE: usize = 10_000;
}

/// `OpenDALBackend` contains a wrapper around an blocking operator of the `OpenDAL` library.
Expand Down Expand Up @@ -217,6 +225,30 @@ impl OpenDALBackend {
Ok(Self { operator })
}

/// Listing options used for repository (and source) listings.
///
/// For B2 this raises the per-request page size from the API default of 100
/// to the documented maximum of 10000.
fn list_options(&self, recursive: bool) -> ListOptions {
ListOptions {
recursive,
limit: self.list_page_size(),
..Default::default()
}
}

/// Backend-specific list page size, if we should override the service default.
fn list_page_size(&self) -> Option<usize> {
let info = self.operator.info();
if !info.capability().list_with_limit {
return None;
}
match info.scheme() {
"b2" => Some(constants::B2_LIST_PAGE_SIZE),
_ => None,
}
}

/// Return a path for the given file type and id.
///
/// # Arguments
Expand Down Expand Up @@ -249,10 +281,7 @@ impl OpenDALBackend {
/// # Errors
/// If listing fails or exclude patterns cannot be compiled
pub fn as_source(self, excludes: &Excludes) -> RusticResult<OpenDALReadSource> {
let list_options = ListOptions {
recursive: true,
..Default::default()
};
let list_options = self.list_options(true);
// openDAL lister may entries in random order; hence we collect and sort them here.
// This also allows to handle listing errors directly
let mut entries: Vec<_> = self
Expand Down Expand Up @@ -348,14 +377,9 @@ impl ReadBackend for OpenDALBackend {
}

let path = tpe.dirname().to_string() + "/";
let list_options = ListOptions {
recursive: true,
..Default::default()
};

let lister = self
.operator
.lister_options(&path, list_options)
.lister_options(&path, self.list_options(true))
.map_err(|err| {
RusticError::with_source(ErrorKind::Backend, "Listing failed for `{type}`", err)
.attach_context("type", tpe.to_string())
Expand Down Expand Up @@ -405,13 +429,9 @@ impl ReadBackend for OpenDALBackend {
}

let path = tpe.dirname().to_string() + "/";
let list_options = ListOptions {
recursive: true,
..Default::default()
};
let lister = self
.operator
.lister_options(&path, list_options)
.lister_options(&path, self.list_options(true))
.map_err(|err| {
RusticError::with_source(ErrorKind::Backend, "Listing failed for `{type}`", err)
.attach_context("type", tpe.to_string())
Expand Down Expand Up @@ -632,6 +652,28 @@ mod tests {
assert!(Throttle::from_str(input).is_err());
}

#[rstest]
#[case("b2", Some(constants::B2_LIST_PAGE_SIZE))]
#[case("s3_aws", None)]
fn list_page_size_matches_scheme(
#[case] fixture: &str,
#[case] expected: Option<usize>,
) -> Result<()> {
#[derive(Deserialize)]
struct TestCase {
path: String,
options: BTreeMap<String, String>,
}

let fixture_path = PathBuf::from(format!("tests/fixtures/opendal/{fixture}.toml"));
let test: TestCase = toml::from_str(&fs::read_to_string(fixture_path)?)?;
let backend = OpenDALBackend::new(test.path, test.options)?;

assert_eq!(backend.list_page_size(), expected);
assert_eq!(backend.list_options(true).limit, expected);
Ok(())
}

#[rstest]
fn new_opendal_backend(
#[files("tests/fixtures/opendal/*.toml")] test_case: PathBuf,
Expand Down
23 changes: 16 additions & 7 deletions crates/core/src/backend/decrypt.rs
Original file line number Diff line number Diff line change
Expand Up @@ -190,17 +190,26 @@ pub trait DecryptReadBackend: ReadBackend + Clone + 'static {
/// If the files could not be read.
fn stream_list<F: RepoFile>(&self, list: Vec<F::Id>, p: &Progress) -> StreamResult<F::Id, F> {
p.set_length(list.len() as u64);
// we use a zero-capacity channel; the loading is typically the bottleneck, not the processing.
let (tx, rx) = bounded(0);
// Index/snapshot files are small; on B2 this is RTT-bound (one GET each).
// Restic uses `connections + GOMAXPROCS`. Keep extra workers so some can
// GET while others decrypt/parse, and buffer so send does not stall IO.
let workers = (rayon::current_num_threads() + 16).clamp(16, 32);
let (tx, rx) = bounded(workers.saturating_mul(2));
let be = self.clone();
let p = p.clone();

spawn(move || {
_ = list.into_par_iter().try_for_each(|id| {
let file = be.get_file::<F>(&id).map(|file| (id, file));
p.inc(1);
tx.send(file).ok() // abort as soon as possible if sending fails, i.e. if the receiver is dropped
});
let work = || {
_ = list.into_par_iter().try_for_each(|id| {
let file = be.get_file::<F>(&id).map(|file| (id, file));
p.inc(1);
tx.send(file).ok()
});
};
match rayon::ThreadPoolBuilder::new().num_threads(workers).build() {
Ok(pool) => pool.install(work),
Err(_) => work(),
}
});
Ok(rx)
}
Expand Down
9 changes: 7 additions & 2 deletions crates/core/src/blob.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,13 @@ pub(super) mod constants {
/// The maximum size of pack-part which is read at once from the backend.
/// (needed to limit the memory size used for large backends)
pub(crate) const LIMIT_PACK_READ: u32 = 40 * 1024 * 1024; // 40 MiB
/// The maximum size of holes which are still read when repacking
pub(crate) const MAX_HOLESIZE: u32 = 256 * 1024; // 256 kiB
/// Maximum unused gap that is still fetched with the surrounding blobs.
///
/// 256 KiB was too small for high-latency object stores (B2): every larger
/// hole became another HTTP range GET, and prune/restore issued those
/// sequentially. 4 MiB is about one RTT of extra download on a ~100 Mbps
/// link, which is cheaper than an extra request.
pub(crate) const MAX_HOLESIZE: u32 = 4 * 1024 * 1024; // 4 MiB
}

/// All [`BlobType`]s which are supported by the repository
Expand Down
28 changes: 21 additions & 7 deletions crates/core/src/blob/tree.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ pub mod rewrite;
use std::{
borrow::Cow,
cmp::Ordering,
collections::{BTreeMap, BTreeSet, BinaryHeap},
collections::{BTreeMap, BinaryHeap, HashSet},
ffi::OsStr,
mem,
path::{Component, Path, PathBuf, Prefix},
Expand All @@ -16,6 +16,7 @@ use crossbeam_channel::{Receiver, Sender, bounded, unbounded};
use derive_setters::Setters;
use ignore::Match;
use ignore::overrides::Override;
use rayon::current_num_threads;
use serde::{Deserialize, Deserializer};
use serde_derive::Serialize;

Expand Down Expand Up @@ -53,8 +54,20 @@ pub enum TreeErrorKind {
pub(crate) type TreeResult<T> = Result<T, TreeErrorKind>;

pub(super) mod constants {
/// The maximum number of trees that are loaded in parallel
pub(super) const MAX_TREE_LOADER: usize = 4;
/// Minimum / maximum tree-loader threads for `TreeStreamerOnce`.
///
/// Four was too few on high-latency backends (B2): prune's "finding used
/// blobs..." walks every unique tree with a pack range GET. Restic uses
/// `connections + GOMAXPROCS` workers. We scale with Rayon (2× CPUs,
/// clamped) so a 4-core box gets 8 loaders, not 4.
pub(super) const MIN_TREE_LOADER: usize = 8;
pub(super) const MAX_TREE_LOADER: usize = 32;
}

fn tree_loader_count() -> usize {
current_num_threads()
.saturating_mul(2)
.clamp(constants::MIN_TREE_LOADER, constants::MAX_TREE_LOADER)
}

pub(crate) type TreeStreamItem = RusticResult<(PathBuf, Tree)>;
Expand Down Expand Up @@ -623,7 +636,7 @@ where
#[derive(Debug)]
pub struct TreeStreamerOnce {
/// The visited tree IDs
visited: BTreeSet<TreeId>,
visited: HashSet<TreeId>,
/// The queue to send tree IDs to
queue_in: Option<Sender<(PathBuf, TreeId, usize)>>,
/// The queue to receive trees from
Expand Down Expand Up @@ -661,10 +674,11 @@ impl TreeStreamerOnce {
) -> RusticResult<Self> {
p.set_length(ids.len() as u64);

let (out_tx, out_rx) = bounded(constants::MAX_TREE_LOADER);
let loaders = tree_loader_count();
let (out_tx, out_rx) = bounded(loaders.saturating_mul(4).max(32));
let (in_tx, in_rx) = unbounded();

for _ in 0..constants::MAX_TREE_LOADER {
for _ in 0..loaders {
let be = be.clone();
let index = index.clone();
let in_rx = in_rx.clone();
Expand All @@ -683,7 +697,7 @@ impl TreeStreamerOnce {

let counter = vec![0; ids.len()];
let mut streamer = Self {
visited: BTreeSet::new(),
visited: HashSet::new(),
queue_in: Some(in_tx),
queue_out: out_rx,
p,
Expand Down
25 changes: 13 additions & 12 deletions crates/core/src/commands/prune.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
/// accessors along with logging macros. Customize as you see fit.
use std::{
cmp::Ordering,
collections::{BTreeMap, BTreeSet},
collections::{BTreeMap, BTreeSet, HashMap},
str::FromStr,
};

Expand Down Expand Up @@ -584,7 +584,7 @@ pub struct PrunePlan {
/// The time the plan was created
time: Zoned,
/// The ids of the blobs which are used
used_ids: BTreeMap<BlobId, u8>,
used_ids: HashMap<BlobId, u8>,
/// The ids of the existing packs
existing_packs: BTreeMap<PackId, u32>,
/// The packs which should be repacked
Expand All @@ -604,7 +604,7 @@ impl PrunePlan {
/// * `existing_packs` - The ids of the existing packs
/// * `index_files` - The index files
fn new(
used_ids: BTreeMap<BlobId, u8>,
used_ids: HashMap<BlobId, u8>,
existing_packs: BTreeMap<PackId, u32>,
index_files: Vec<(IndexId, IndexFile)>,
) -> Self {
Expand Down Expand Up @@ -1416,14 +1416,15 @@ pub(crate) fn prune_repository<S: Open>(
})
.collect();

// TODO: repack in parallel
for blobs in blob_chunks {
// Range-GETs for holes in the same pack run in parallel. The
// packer already serializes writes via its channel / lock.
blob_chunks.into_par_iter().try_for_each(|blobs| {
if opts.fast_repack {
repacker.copy_fast(blobs, &p)?;
repacker.copy_fast(blobs, &p)
} else {
repacker.copy(blobs, &p)?;
repacker.copy(blobs, &p)
}
}
})?;
Ok(())
})?;
_ = tree_repacker.finalize()?;
Expand Down Expand Up @@ -1491,8 +1492,8 @@ impl PackInfo {
/// # Arguments
///
/// * `pack` - The `PrunePack` to create the `PackInfo` from
/// * `used_ids` - The `BTreeMap` of used ids
fn from_pack(pack: &PrunePack, used_ids: &mut BTreeMap<BlobId, u8>) -> Self {
/// * `used_ids` - The map of used ids
fn from_pack(pack: &PrunePack, used_ids: &mut HashMap<BlobId, u8>) -> Self {
let mut pi = Self {
blob_type: pack.blob_type,
used_blobs: 0,
Expand Down Expand Up @@ -1584,7 +1585,7 @@ fn find_used_blobs<S>(
be: &impl DecryptReadBackend,
index: &impl ReadGlobalIndex,
ignore_snaps: &[SnapshotId],
) -> RusticResult<BTreeMap<BlobId, u8>> {
) -> RusticResult<HashMap<BlobId, u8>> {
let ignore_snaps: BTreeSet<_> = ignore_snaps.iter().collect();

let p = repo.progress_counter("reading snapshots...");
Expand All @@ -1601,7 +1602,7 @@ fn find_used_blobs<S>(
.try_collect()?;
p.finish();

let mut ids: BTreeMap<_, _> = snap_trees
let mut ids: HashMap<_, _> = snap_trees
.iter()
.map(|id| (BlobId::from(**id), 0))
.collect();
Expand Down
Loading