Remote files open in a few requests; ObjectHeader::parse back to speed; last conformance mismatches resolved (602/697) #21
@@ -89,6 +89,21 @@ pub trait Storage {
|
|||||||
fn as_contiguous(&self) -> Option<&[u8]> {
|
fn as_contiguous(&self) -> Option<&[u8]> {
|
||||||
None
|
None
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A hint that `[offset, offset + len)` is about to be read by the
|
||||||
|
/// same operation: a parser that has just learnt where the next
|
||||||
|
/// structures are (a node's children, a structure's body) says so
|
||||||
|
/// before it reads them one at a time. Nothing is read and nothing
|
||||||
|
/// fails. The default ignores it, as does every backend that reads
|
||||||
|
/// when asked; the browser's restartable reader, which fetches over
|
||||||
|
/// the network between attempts, fetches hinted bytes along with the
|
||||||
|
/// bytes an attempt actually missed, so structures a parser only
|
||||||
|
/// reaches after a miss arrive in the same round trip. Results never
|
||||||
|
/// depend on hints.
|
||||||
|
#[inline]
|
||||||
|
fn hint(&self, offset: u64, len: usize) {
|
||||||
|
let _ = (offset, len);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Storage for [u8] {
|
impl Storage for [u8] {
|
||||||
@@ -148,6 +163,11 @@ impl<T: Storage + ?Sized> Storage for &T {
|
|||||||
fn as_contiguous(&self) -> Option<&[u8]> {
|
fn as_contiguous(&self) -> Option<&[u8]> {
|
||||||
(**self).as_contiguous()
|
(**self).as_contiguous()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[inline]
|
||||||
|
fn hint(&self, offset: u64, len: usize) {
|
||||||
|
(**self).hint(offset, len)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<T: Storage + ?Sized> Storage for Box<T> {
|
impl<T: Storage + ?Sized> Storage for Box<T> {
|
||||||
@@ -170,6 +190,11 @@ impl<T: Storage + ?Sized> Storage for Box<T> {
|
|||||||
fn as_contiguous(&self) -> Option<&[u8]> {
|
fn as_contiguous(&self) -> Option<&[u8]> {
|
||||||
(**self).as_contiguous()
|
(**self).as_contiguous()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[inline]
|
||||||
|
fn hint(&self, offset: u64, len: usize) {
|
||||||
|
(**self).hint(offset, len)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(feature = "std")]
|
#[cfg(feature = "std")]
|
||||||
@@ -193,6 +218,11 @@ impl<T: Storage + ?Sized> Storage for std::sync::Arc<T> {
|
|||||||
fn as_contiguous(&self) -> Option<&[u8]> {
|
fn as_contiguous(&self) -> Option<&[u8]> {
|
||||||
(**self).as_contiguous()
|
(**self).as_contiguous()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[inline]
|
||||||
|
fn hint(&self, offset: u64, len: usize) {
|
||||||
|
(**self).hint(offset, len)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// `storage.len()` as the `usize` the parsers' end-of-file errors report
|
/// `storage.len()` as the `usize` the parsers' end-of-file errors report
|
||||||
|
|||||||
@@ -25,6 +25,14 @@
|
|||||||
//! pass per block it needs. In practice it is one pass per *wave* of
|
//! pass per block it needs. In practice it is one pass per *wave* of
|
||||||
//! misses: a chunked read asks for all the chunks of a batch at once.
|
//! misses: a chunked read asks for all the chunks of a batch at once.
|
||||||
//!
|
//!
|
||||||
|
//! Parsers also say what they are about to read ([`Storage::hint`]: a
|
||||||
|
//! node's children, a structure's body, a listed group's child headers).
|
||||||
|
//! A pass that misses fetches the hinted blocks it lacks too, as far as the
|
||||||
|
//! operation's fetch budget allows, so what the parser would only have
|
||||||
|
//! reached on the next pass arrives in the same round trip; a pass that
|
||||||
|
//! misses nothing ignores its hints, so they never add a round trip, and
|
||||||
|
//! results never depend on them.
|
||||||
|
//!
|
||||||
//! Blocks are kept in an LRU cache with a byte budget, trimmed only when no
|
//! Blocks are kept in an LRU cache with a byte budget, trimmed only when no
|
||||||
//! operation is in flight. Blocks fetched for bulk reads (raw data: a
|
//! operation is in flight. Blocks fetched for bulk reads (raw data: a
|
||||||
//! `read_ranges` call, or a read longer than a block) go first, so reading
|
//! `read_ranges` call, or a read longer than a block) go first, so reading
|
||||||
@@ -50,6 +58,14 @@ pub const DEFAULT_MAX_FETCH: u64 = 512 << 20;
|
|||||||
/// the caller of [`LazyStorage::attempt`]: a pass that missed is re-run.
|
/// the caller of [`LazyStorage::attempt`]: a pass that missed is re-run.
|
||||||
pub const NEED_BYTES: &str = "bytes not fetched yet (restartable read)";
|
pub const NEED_BYTES: &str = "bytes not fetched yet (restartable read)";
|
||||||
|
|
||||||
|
/// Most blocks one pass records as hinted (see [`Storage::hint`]).
|
||||||
|
const MAX_HINTED: usize = 1 << 16;
|
||||||
|
|
||||||
|
/// Total length of `ranges`.
|
||||||
|
fn ranges_len(ranges: &[Range<u64>]) -> u64 {
|
||||||
|
ranges.iter().map(|r| r.end - r.start).sum()
|
||||||
|
}
|
||||||
|
|
||||||
/// Settings of a [`LazyStorage`].
|
/// Settings of a [`LazyStorage`].
|
||||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
pub struct LazyConfig {
|
pub struct LazyConfig {
|
||||||
@@ -94,6 +110,9 @@ pub struct LazyStats {
|
|||||||
pub bytes_fetched: u64,
|
pub bytes_fetched: u64,
|
||||||
/// Blocks evicted to stay within the budget.
|
/// Blocks evicted to stay within the budget.
|
||||||
pub evictions: u64,
|
pub evictions: u64,
|
||||||
|
/// Blocks asked for because a parser hinted it would read them
|
||||||
|
/// ([`Storage::hint`]), not because a pass missed them.
|
||||||
|
pub hinted_blocks: u64,
|
||||||
/// Bytes cached now.
|
/// Bytes cached now.
|
||||||
pub cached_bytes: u64,
|
pub cached_bytes: u64,
|
||||||
}
|
}
|
||||||
@@ -125,6 +144,9 @@ struct State {
|
|||||||
/// Blocks the current pass missed, and whether a small read wanted
|
/// Blocks the current pass missed, and whether a small read wanted
|
||||||
/// them (metadata).
|
/// them (metadata).
|
||||||
missing: HashMap<u64, bool>,
|
missing: HashMap<u64, bool>,
|
||||||
|
/// Blocks the current pass was told it is about to read and that are
|
||||||
|
/// not cached ([`Storage::hint`]), in file order.
|
||||||
|
hinted: BTreeSet<u64>,
|
||||||
/// Blocks a bulk read missed that have not been supplied yet: kept
|
/// Blocks a bulk read missed that have not been supplied yet: kept
|
||||||
/// as bulk when they arrive.
|
/// as bulk when they arrive.
|
||||||
bulk_pending: BTreeSet<u64>,
|
bulk_pending: BTreeSet<u64>,
|
||||||
@@ -155,6 +177,19 @@ pub struct Operation<'a> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl Operation<'_> {
|
impl Operation<'_> {
|
||||||
|
/// One pass of `f`, as [`LazyStorage::attempt`], but a pass that misses
|
||||||
|
/// also asks for the blocks it was hinted it would read, as many as fit
|
||||||
|
/// in what is left of the operation's budget ([`LazyConfig::max_fetch`])
|
||||||
|
/// after the blocks it missed.
|
||||||
|
pub fn attempt<T>(&self, f: impl FnOnce() -> T) -> Step<T> {
|
||||||
|
let left = self
|
||||||
|
.storage
|
||||||
|
.config
|
||||||
|
.max_fetch
|
||||||
|
.saturating_sub(self.fetched.get());
|
||||||
|
self.storage.attempt_within(f, left)
|
||||||
|
}
|
||||||
|
|
||||||
/// Count `ranges` against the operation's budget
|
/// Count `ranges` against the operation's budget
|
||||||
/// ([`LazyConfig::max_fetch`]) before they are fetched: an error, and
|
/// ([`LazyConfig::max_fetch`]) before they are fetched: an error, and
|
||||||
/// nothing counted, if they would take it past the budget.
|
/// nothing counted, if they would take it past the budget.
|
||||||
@@ -225,20 +260,70 @@ impl LazyStorage {
|
|||||||
/// Run one pass of `f` over this storage. `Done` when `f` read nothing
|
/// Run one pass of `f` over this storage. `Done` when `f` read nothing
|
||||||
/// that is missing; otherwise `Need` with the ranges to fetch, and `f`'s
|
/// that is missing; otherwise `Need` with the ranges to fetch, and `f`'s
|
||||||
/// result is dropped (it may be an error caused by the miss, or a
|
/// result is dropped (it may be an error caused by the miss, or a
|
||||||
/// result built around one).
|
/// result built around one). Hints are not followed; see
|
||||||
|
/// [`Operation::attempt`].
|
||||||
pub fn attempt<T>(&self, f: impl FnOnce() -> T) -> Step<T> {
|
pub fn attempt<T>(&self, f: impl FnOnce() -> T) -> Step<T> {
|
||||||
|
self.attempt_within(f, 0)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// [`attempt`](Self::attempt), adding to a pass that misses the hinted
|
||||||
|
/// blocks it lacks while everything asked for stays within `budget`
|
||||||
|
/// bytes.
|
||||||
|
fn attempt_within<T>(&self, f: impl FnOnce() -> T, budget: u64) -> Step<T> {
|
||||||
{
|
{
|
||||||
let mut st = lock(&self.state);
|
let mut st = lock(&self.state);
|
||||||
st.missing.clear();
|
st.missing.clear();
|
||||||
|
st.hinted.clear();
|
||||||
st.stats.passes += 1;
|
st.stats.passes += 1;
|
||||||
}
|
}
|
||||||
let out = f();
|
let out = f();
|
||||||
let missing = std::mem::take(&mut lock(&self.state).missing);
|
let (missing, hinted) = {
|
||||||
|
let mut st = lock(&self.state);
|
||||||
|
(
|
||||||
|
std::mem::take(&mut st.missing),
|
||||||
|
std::mem::take(&mut st.hinted),
|
||||||
|
)
|
||||||
|
};
|
||||||
if missing.is_empty() {
|
if missing.is_empty() {
|
||||||
return Step::Done(out);
|
return Step::Done(out);
|
||||||
}
|
}
|
||||||
drop(out);
|
drop(out);
|
||||||
Step::Need(self.runs(missing))
|
let need = self.runs(&missing);
|
||||||
|
if hinted.is_empty() || budget == 0 {
|
||||||
|
return Step::Need(need);
|
||||||
|
}
|
||||||
|
// The hinted blocks still missing, in file order, while they fit.
|
||||||
|
let bs = self.config.block_size;
|
||||||
|
let mut total = ranges_len(&need);
|
||||||
|
let mut wanted = missing.clone();
|
||||||
|
{
|
||||||
|
let st = lock(&self.state);
|
||||||
|
for i in hinted {
|
||||||
|
if wanted.contains_key(&i) || st.blocks.contains_key(&i) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
let n = (self.len - i * bs).min(bs);
|
||||||
|
if total + n > budget {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
total += n;
|
||||||
|
wanted.insert(i, true);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if wanted.len() == missing.len() {
|
||||||
|
return Step::Need(need);
|
||||||
|
}
|
||||||
|
let with_hints = self.runs(&wanted);
|
||||||
|
// Filling holes between runs can add blocks: never let the hints
|
||||||
|
// take the pass past the budget.
|
||||||
|
if ranges_len(&with_hints) > budget {
|
||||||
|
return Step::Need(need);
|
||||||
|
}
|
||||||
|
{
|
||||||
|
let mut st = lock(&self.state);
|
||||||
|
st.stats.hinted_blocks += (wanted.len() - missing.len()) as u64;
|
||||||
|
}
|
||||||
|
Step::Need(with_hints)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// The bytes of the file at `offset`, fetched for a range a pass asked
|
/// The bytes of the file at `offset`, fetched for a range a pass asked
|
||||||
@@ -309,7 +394,7 @@ impl LazyStorage {
|
|||||||
) -> Result<T, String> {
|
) -> Result<T, String> {
|
||||||
let op = self.operation();
|
let op = self.operation();
|
||||||
loop {
|
loop {
|
||||||
match self.attempt(&mut f) {
|
match op.attempt(&mut f) {
|
||||||
Step::Done(v) => return Ok(v),
|
Step::Done(v) => return Ok(v),
|
||||||
Step::Need(ranges) => {
|
Step::Need(ranges) => {
|
||||||
op.charge(&ranges)?;
|
op.charge(&ranges)?;
|
||||||
@@ -350,14 +435,14 @@ impl LazyStorage {
|
|||||||
/// blocks, a one-block hole between two runs filled so they merge
|
/// blocks, a one-block hole between two runs filled so they merge
|
||||||
/// (unless the hole is cached: it would be fetched again), each at
|
/// (unless the hole is cached: it would be fetched again), each at
|
||||||
/// most `max_request` long.
|
/// most `max_request` long.
|
||||||
fn runs(&self, missing: HashMap<u64, bool>) -> Vec<Range<u64>> {
|
fn runs(&self, missing: &HashMap<u64, bool>) -> Vec<Range<u64>> {
|
||||||
let bs = self.config.block_size;
|
let bs = self.config.block_size;
|
||||||
let mut wanted: Vec<u64> = missing.keys().copied().collect();
|
let mut wanted: Vec<u64> = missing.keys().copied().collect();
|
||||||
wanted.sort_unstable();
|
wanted.sort_unstable();
|
||||||
let mut st = lock(&self.state);
|
let mut st = lock(&self.state);
|
||||||
// Remember which blocks only bulk reads asked for: they are kept
|
// Remember which blocks only bulk reads asked for: they are kept
|
||||||
// as bulk once supplied.
|
// as bulk once supplied.
|
||||||
for (&i, &metadata) in &missing {
|
for (&i, &metadata) in missing {
|
||||||
if metadata {
|
if metadata {
|
||||||
st.bulk_pending.remove(&i);
|
st.bulk_pending.remove(&i);
|
||||||
} else {
|
} else {
|
||||||
@@ -499,6 +584,25 @@ impl Storage for LazyStorage {
|
|||||||
self.len
|
self.len
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn hint(&self, offset: u64, len: usize) {
|
||||||
|
// A hint is about a structure, not bulk data: at most a block or
|
||||||
|
// 1 MiB of it is followed, and at most `MAX_HINTED` blocks a pass,
|
||||||
|
// whatever a hostile file makes a parser hint.
|
||||||
|
let len = (len as u64).min(self.config.block_size.max(1 << 20));
|
||||||
|
let Some(span) = self.span(offset, len) else {
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
let mut st = lock(&self.state);
|
||||||
|
for i in span {
|
||||||
|
if st.hinted.len() >= MAX_HINTED {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if !st.blocks.contains_key(&i) {
|
||||||
|
st.hinted.insert(i);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fn read_ranges(&self, ranges: &[Range<u64>]) -> Result<Vec<Cow<'_, [u8]>>, FormatError> {
|
fn read_ranges(&self, ranges: &[Range<u64>]) -> Result<Vec<Cow<'_, [u8]>>, FormatError> {
|
||||||
let mut spans = Vec::with_capacity(ranges.len());
|
let mut spans = Vec::with_capacity(ranges.len());
|
||||||
let mut total = 0u64;
|
let mut total = 0u64;
|
||||||
@@ -797,6 +901,89 @@ mod tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn hinted_blocks_come_with_a_miss_and_never_alone() {
|
||||||
|
let data = file(16 * 1024);
|
||||||
|
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
|
||||||
|
serve(&s, &data, &[3072..4096]);
|
||||||
|
let op = s.operation();
|
||||||
|
// Hints alone: the pass is done, nothing is fetched.
|
||||||
|
let step = op.attempt(|| {
|
||||||
|
s.hint(8 * 1024, 100);
|
||||||
|
owned(s.read_at(3072, 8))
|
||||||
|
});
|
||||||
|
assert!(matches!(step, Step::Done(Ok(_))), "{step:?}");
|
||||||
|
// With a miss, the hinted blocks not cached come too (block 3 is
|
||||||
|
// cached; 10..=11 is one run, 14 another).
|
||||||
|
let pass = || {
|
||||||
|
s.hint(3072, 10);
|
||||||
|
s.hint(10 * 1024 + 1000, 100);
|
||||||
|
s.hint(14 * 1024, 1);
|
||||||
|
owned(s.read_at(0, 8))
|
||||||
|
};
|
||||||
|
let Step::Need(need) = op.attempt(pass) else {
|
||||||
|
panic!("block 0 is missing");
|
||||||
|
};
|
||||||
|
assert_eq!(
|
||||||
|
need,
|
||||||
|
vec![0..1024, 10 * 1024..12 * 1024, 14 * 1024..15 * 1024]
|
||||||
|
);
|
||||||
|
// A plain attempt does not follow hints.
|
||||||
|
let Step::Need(plain) = s.attempt(pass) else {
|
||||||
|
panic!("block 0 is missing");
|
||||||
|
};
|
||||||
|
assert_eq!(plain, vec![0..1024]);
|
||||||
|
serve(&s, &data, &need);
|
||||||
|
let Step::Done(got) = op.attempt(pass) else {
|
||||||
|
panic!("everything was supplied");
|
||||||
|
};
|
||||||
|
assert_eq!(got.unwrap(), &data[..8]);
|
||||||
|
assert_eq!(s.stats().hinted_blocks, 3);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn hints_stay_within_the_fetch_budget() {
|
||||||
|
// A budget of 3 blocks: the missed block and the first two hinted
|
||||||
|
// ones fit, the rest are left out; a hint never makes a call fail.
|
||||||
|
// (Hinted blocks two apart would be merged with the hole between
|
||||||
|
// them, which would not fit: then no hint is followed.)
|
||||||
|
let data = file(64 * 1024);
|
||||||
|
let mut c = config(1024, 1 << 20);
|
||||||
|
c.max_fetch = 3 * 1024;
|
||||||
|
let s = LazyStorage::new(data.len() as u64, c);
|
||||||
|
let got = s
|
||||||
|
.run_blocking(
|
||||||
|
|| {
|
||||||
|
for i in 0..20 {
|
||||||
|
s.hint(20 * 1024 + i * 3072, 1);
|
||||||
|
}
|
||||||
|
owned(s.read_at(0, 8))
|
||||||
|
},
|
||||||
|
|r| Ok(data[r.start as usize..r.end as usize].to_vec()),
|
||||||
|
)
|
||||||
|
.unwrap()
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(got, &data[..8]);
|
||||||
|
let st = s.stats();
|
||||||
|
assert_eq!(
|
||||||
|
(st.bytes_fetched, st.hinted_blocks),
|
||||||
|
(3 * 1024, 2),
|
||||||
|
"{st:?}"
|
||||||
|
);
|
||||||
|
// A hint longer than the file, or past its end, is harmless.
|
||||||
|
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
|
||||||
|
s.run_blocking(
|
||||||
|
|| {
|
||||||
|
s.hint(0, usize::MAX);
|
||||||
|
s.hint(u64::MAX, 10);
|
||||||
|
owned(s.read_at(0, 8))
|
||||||
|
},
|
||||||
|
|r| Ok(data[r.start as usize..r.end as usize].to_vec()),
|
||||||
|
)
|
||||||
|
.unwrap()
|
||||||
|
.unwrap();
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn a_failed_fetch_is_an_error_not_data() {
|
fn a_failed_fetch_is_an_error_not_data() {
|
||||||
let data = file(4096);
|
let data = file(4096);
|
||||||
|
|||||||
@@ -391,7 +391,7 @@ impl Http {
|
|||||||
) -> Result<T, JsError> {
|
) -> Result<T, JsError> {
|
||||||
let op = storage.operation();
|
let op = storage.operation();
|
||||||
loop {
|
loop {
|
||||||
match storage.attempt(&mut f) {
|
match op.attempt(&mut f) {
|
||||||
Step::Done(v) => return Ok(v),
|
Step::Done(v) => return Ok(v),
|
||||||
Step::Need(ranges) => {
|
Step::Need(ranges) => {
|
||||||
op.charge(&ranges).map_err(js_err)?;
|
op.charge(&ranges).map_err(js_err)?;
|
||||||
|
|||||||
@@ -372,6 +372,21 @@ impl Storage for FileData {
|
|||||||
fn as_contiguous(&self) -> Option<&[u8]> {
|
fn as_contiguous(&self) -> Option<&[u8]> {
|
||||||
self.contiguous()
|
self.contiguous()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn hint(&self, offset: u64, len: usize) {
|
||||||
|
if self.contiguous().is_some() {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if let Backing::Storage(s) = &self.backing {
|
||||||
|
// Within the HDF5 data, as a read would be clamped.
|
||||||
|
let size = Storage::len(self);
|
||||||
|
let start = offset.min(size);
|
||||||
|
let len = usize::try_from(size - start).map_or(len, |avail| avail.min(len));
|
||||||
|
if len > 0 {
|
||||||
|
s.hint(self.base + start, len);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// `bytes`, a backend's answer to a read of `len` bytes, without anything
|
/// `bytes`, a backend's answer to a read of `len` bytes, without anything
|
||||||
|
|||||||
Reference in New Issue
Block a user