diff --git a/lib/common/common/benches/universal_io.rs b/lib/common/common/benches/universal_io.rs index a0b99d2ca0..4a0f99d9fd 100644 --- a/lib/common/common/benches/universal_io.rs +++ b/lib/common/common/benches/universal_io.rs @@ -110,12 +110,14 @@ fn read_benches>( group.bench_function("read_batch_small_random", |b| { b.iter(|| { let mut sum = 0u64; - let ranges = (0..8).map(|_| ReadRange { - byte_offset: rng.random_range(0..len) * size_of::() as u64, - length: 1, - }); + let ranges = (0..8) + .map(|_| ReadRange { + byte_offset: rng.random_range(0..len) * size_of::() as u64, + length: 1, + }) + .map(|range| ((), range)); storage - .read_batch::(ranges, |_, chunk| { + .read_batch::(ranges, |(), chunk| { for &item in bytemuck::cast_slice::(chunk) { sum = sum.wrapping_add(item); } @@ -131,12 +133,14 @@ fn read_benches>( b.iter(|| { let mut sum = 0u64; let start = rng.random_range(0..len - 8) * size_of::() as u64; - let ranges = (0..8).map(move |i| ReadRange { - byte_offset: start + i * size_of::() as u64, - length: 1, - }); + let ranges = (0..8) + .map(move |i| ReadRange { + byte_offset: start + i * size_of::() as u64, + length: 1, + }) + .map(|range| ((), range)); storage - .read_batch::(ranges, |_, chunk| { + .read_batch::(ranges, |(), chunk| { for &item in bytemuck::cast_slice::(chunk) { sum = sum.wrapping_add(item); } @@ -152,7 +156,7 @@ fn read_benches>( let read_batch_full = || { let mut sum = 0u64; storage - .read_batch::(ranges_full_file::(), |_, chunk| { + .read_batch::(ranges_full_file::(), |(), chunk| { for &item in bytemuck::cast_slice::(chunk) { sum = sum.wrapping_add(item); } @@ -175,12 +179,14 @@ fn time_it(f: F) -> Duration { start.elapsed() } -fn ranges_full_file() -> impl Iterator { +fn ranges_full_file() -> impl Iterator { let len = FILE_SIZE_BYTES / size_of::() as u64; - (0..len).map(move |i| ReadRange { - byte_offset: i * size_of::() as u64, - length: 1, - }) + (0..len) + .map(move |i| ReadRange { + byte_offset: i * size_of::() as u64, + length: 1, + }) + .map(|range| ((), range)) } fn make_random_file() -> PathBuf { diff --git a/lib/common/common/src/universal_io/disk_cache/mod.rs b/lib/common/common/src/universal_io/disk_cache/mod.rs index 570d51e728..b3b027c768 100644 --- a/lib/common/common/src/universal_io/disk_cache/mod.rs +++ b/lib/common/common/src/universal_io/disk_cache/mod.rs @@ -96,14 +96,14 @@ impl UniversalRead for CachedSlice { Ok(self.get_range(range)?) } - fn read_batch( - &self, - ranges: impl IntoIterator, - mut callback: impl FnMut(usize, &[T]) -> Result<()>, + fn read_batch<'a, P: AccessPattern, Meta: 'a>( + &'a self, + ranges: impl IntoIterator, + mut callback: impl FnMut(Meta, &[T]) -> Result<()>, ) -> Result<()> { - for (i, range) in ranges.into_iter().enumerate() { + for (meta, range) in ranges { let data = self.read::

(range)?; - callback(i, &data)?; + callback(meta, &data)?; } Ok(()) diff --git a/lib/common/common/src/universal_io/io_uring/mod.rs b/lib/common/common/src/universal_io/io_uring/mod.rs index 7f406a64e7..f449d8c7c8 100644 --- a/lib/common/common/src/universal_io/io_uring/mod.rs +++ b/lib/common/common/src/universal_io/io_uring/mod.rs @@ -91,48 +91,62 @@ impl UniversalRead for IoUringFile { Ok(Cow::Owned(items)) } - fn read_batch( - &self, - ranges: impl IntoIterator, - mut callback: impl FnMut(usize, &[T]) -> Result<()>, + fn read_batch<'a, P: AccessPattern, Meta: 'a>( + &'a self, + ranges: impl IntoIterator, + mut callback: impl FnMut(Meta, &[T]) -> Result<()>, ) -> Result<()> { - for record in self.read_iter::

(ranges) { - let (idx, items) = record?; - callback(idx, &items)?; + for record in self.read_iter::(ranges) { + let (meta, items) = record?; + callback(meta, &items)?; } Ok(()) } - fn read_iter( + fn read_iter( &self, - ranges: impl IntoIterator, - ) -> impl Iterator)>> { - match IoUringReadIter::new(self, ranges.into_iter()) { - Ok(iter) => itertools::Either::Left(iter), + ranges: impl IntoIterator, + ) -> impl Iterator)>> { + let fd = self.fd(); + let direct_io = self.direct_io; + let ranges = ranges + .into_iter() + .map(move |(meta, range)| (meta, fd, direct_io, range)); + match IoUringReadIter::new(ranges) { + Ok(iter) => itertools::Either::Left( + iter.map(|result| result.map(|(meta, items)| (meta, Cow::Owned(items)))), + ), Err(err) => itertools::Either::Right(iter::once(Err(err))), } } - fn read_multi( - files: &[Self], - reads: impl IntoIterator, - mut callback: impl FnMut(usize, FileIndex, &[T]) -> Result<()>, - ) -> Result<()> { - for record in Self::read_multi_iter::

(files, reads) { - let (idx, file_idx, items) = record?; - callback(idx, file_idx, &items)?; + fn read_multi<'a, P: AccessPattern, Meta: 'a>( + reads: impl IntoIterator, + mut callback: impl FnMut(Meta, &[T]) -> Result<()>, + ) -> Result<()> + where + Self: 'a, + { + for record in Self::read_multi_iter::<'a, P, Meta>(reads) { + let (meta, items) = record?; + callback(meta, &items)?; } Ok(()) } - fn read_multi_iter( - files: &[Self], - reads: impl IntoIterator, - ) -> impl Iterator)>> { - match IoUringReadMultiIter::new(files, reads.into_iter()) { - Ok(iter) => itertools::Either::Left(iter), + fn read_multi_iter<'a, P: AccessPattern, Meta>( + reads: impl IntoIterator, + ) -> impl Iterator)>> { + let ranges = reads + .into_iter() + .map(|(meta, file, range)| (meta, file.fd(), file.direct_io, range)); + + match IoUringReadIter::new(ranges) { + Ok(iter) => itertools::Either::Left( + iter.map(|result| result.map(|(meta, items)| (meta, Cow::Owned(items)))), + ), Err(err) => itertools::Either::Right(iter::once(Err(err))), } } @@ -180,15 +194,15 @@ impl UniversalWrite for IoUringFile { items: impl IntoIterator, ) -> Result<()> { let mut rt = IoUringRuntime::new()?; - let mut items = items.into_iter().enumerate().peekable(); + let mut items = items.into_iter().peekable(); while items.peek().is_some() || rt.in_progress > 0 { rt.enqueue_while(|state| { - let Some((id, (byte_offset, items))) = items.next() else { + let Some((byte_offset, items)) = items.next() else { return Ok(None); }; - let entry = state.write(id as _, self.fd(), byte_offset, items); + let entry = state.write((), self.fd(), byte_offset, items); Ok(Some(entry)) })?; @@ -208,11 +222,11 @@ impl UniversalWrite for IoUringFile { writes: impl IntoIterator, ) -> Result<()> { let mut rt = IoUringRuntime::new()?; - let mut writes = writes.into_iter().enumerate().peekable(); + let mut writes = writes.into_iter().peekable(); while writes.peek().is_some() || rt.in_progress > 0 { rt.enqueue_while(|state| { - let Some((id, (file_index, byte_offset, items))) = writes.next() else { + let Some((file_index, byte_offset, items)) = writes.next() else { return Ok(None); }; @@ -223,7 +237,7 @@ impl UniversalWrite for IoUringFile { } })?; - let entry = state.write(id as _, file.fd(), byte_offset, items); + let entry = state.write((), file.fd(), byte_offset, items); Ok(Some(entry)) })?; diff --git a/lib/common/common/src/universal_io/io_uring/read_iter.rs b/lib/common/common/src/universal_io/io_uring/read_iter.rs index 110a6c48e5..0016eacb36 100644 --- a/lib/common/common/src/universal_io/io_uring/read_iter.rs +++ b/lib/common/common/src/universal_io/io_uring/read_iter.rs @@ -1,25 +1,24 @@ -use std::borrow::Cow; use std::iter; use std::marker::PhantomData; +use ::io_uring::types::Fd; + use super::*; -pub struct IoUringReadIter<'a, T: 'static, I: Iterator> { - file: &'a IoUringFile, - ranges: iter::Peekable>, - runtime: IoUringRuntime<'static, T>, +pub struct IoUringReadIter { + ranges: iter::Peekable, + runtime: IoUringRuntime<'static, T, Meta>, _phantom: PhantomData<*const ()>, // `!Send + !Sync` } -impl<'a, T, I> IoUringReadIter<'a, T, I> +impl IoUringReadIter where T: bytemuck::Pod, - I: Iterator, + I: Iterator, { - pub fn new(file: &'a IoUringFile, ranges: I) -> Result { + pub fn new(ranges: I) -> Result { let iter = Self { - file, - ranges: ranges.enumerate().peekable(), + ranges: ranges.peekable(), runtime: IoUringRuntime::new()?, _phantom: PhantomData, }; @@ -27,16 +26,15 @@ where Ok(iter) } - fn next_impl(&mut self) -> Result)>> { + fn next_impl(&mut self) -> Result)>> { if self.runtime.completion_is_empty() && (self.ranges.peek().is_some() || self.runtime.in_progress > 0) { self.runtime.enqueue_while(|state| { - let Some((id, range)) = self.ranges.next() else { + let Some((meta, fd, direct_io, range)) = self.ranges.next() else { return Ok(None); }; - - let entry = state.read(id as _, self.file.fd(), range, 0, self.file.direct_io); + let entry = state.read(meta, fd, range, direct_io); Ok(Some(entry)) })?; @@ -48,99 +46,18 @@ where .completed() .next() .transpose()? - .map(|(id, resp)| { - let id = id as _; - let (_, items) = resp.expect_read(); - - (id, Cow::from(items)) - }); + .map(|(meta, resp)| (meta, resp.expect_read())); Ok(next) } } -impl<'a, T, I> Iterator for IoUringReadIter<'a, T, I> +impl Iterator for IoUringReadIter where T: bytemuck::Pod, - I: Iterator, + I: Iterator, { - type Item = Result<(usize, Cow<'a, [T]>)>; - - fn next(&mut self) -> Option { - self.next_impl().transpose() - } -} - -pub struct IoUringReadMultiIter<'a, T: 'static, I: Iterator> { - files: &'a [IoUringFile], - ranges: iter::Peekable>, - runtime: IoUringRuntime<'static, T>, - _phantom: PhantomData<*const ()>, // `!Send + !Sync` -} - -impl<'a, T, I> IoUringReadMultiIter<'a, T, I> -where - T: bytemuck::Pod, - I: Iterator, -{ - pub fn new(files: &'a [IoUringFile], ranges: I) -> Result { - let iter = Self { - files, - ranges: ranges.enumerate().peekable(), - runtime: IoUringRuntime::new()?, - _phantom: PhantomData, - }; - - Ok(iter) - } - - #[expect(clippy::type_complexity)] - fn next_impl(&mut self) -> Result)>> { - if self.runtime.completion_is_empty() - && (self.ranges.peek().is_some() || self.runtime.in_progress > 0) - { - self.runtime.enqueue_while(|state| { - let Some((id, (file_index, range))) = self.ranges.next() else { - return Ok(None); - }; - - let file = - self.files - .get(file_index) - .ok_or(UniversalIoError::InvalidFileIndex { - file_index, - files: self.files.len(), - })?; - - let entry = state.read(id as _, file.fd(), range, file_index, file.direct_io); - Ok(Some(entry)) - })?; - - self.runtime.submit_and_wait(1)?; - } - - let next = self - .runtime - .completed() - .next() - .transpose()? - .map(|(id, resp)| { - let id = id as _; - let (file_index, items) = resp.expect_read(); - - (id, file_index, Cow::from(items)) - }); - - Ok(next) - } -} - -impl<'a, T, I> Iterator for IoUringReadMultiIter<'a, T, I> -where - T: bytemuck::Pod, - I: Iterator, -{ - type Item = Result<(usize, FileIndex, Cow<'a, [T]>)>; + type Item = Result<(Meta, Vec)>; fn next(&mut self) -> Option { self.next_impl().transpose() diff --git a/lib/common/common/src/universal_io/io_uring/runtime.rs b/lib/common/common/src/universal_io/io_uring/runtime.rs index fe162c3e4b..f87bc19295 100644 --- a/lib/common/common/src/universal_io/io_uring/runtime.rs +++ b/lib/common/common/src/universal_io/io_uring/runtime.rs @@ -8,13 +8,13 @@ use slab::Slab; use super::*; use crate::maybe_uninit; -pub struct IoUringRuntime<'data, T> { +pub struct IoUringRuntime<'data, T, Meta = u64> { io_uring: IoUringGuard, - state: IoUringState<'data, T>, + state: IoUringState<'data, T, Meta>, pub in_progress: usize, } -impl<'data, T> IoUringRuntime<'data, T> { +impl<'data, T, Meta> IoUringRuntime<'data, T, Meta> { pub fn new() -> Result { let mut io_uring = pool::get_io_uring()?; let capacity = io_uring.submission().capacity(); @@ -31,7 +31,7 @@ impl<'data, T> IoUringRuntime<'data, T> { /// or the queue is full. pub fn enqueue_while(&mut self, mut entries: F) -> Result<()> where - F: FnMut(&mut IoUringState<'data, T>) -> Result>, + F: FnMut(&mut IoUringState<'data, T, Meta>) -> Result>, { let mut squeue = self.io_uring.submission(); @@ -97,7 +97,7 @@ impl<'data, T> IoUringRuntime<'data, T> { Ok(()) } - pub fn completed(&mut self) -> impl Iterator)>> { + pub fn completed(&mut self) -> impl Iterator)>> { self.io_uring.completion().map(|entry| { self.in_progress -= 1; @@ -114,13 +114,13 @@ impl<'data, T> IoUringRuntime<'data, T> { } let length = result as _; - let (id, resp) = self.state.finalize(slot, length)?; - Ok((id, resp)) + let (meta, resp) = self.state.finalize(slot, length)?; + Ok((meta, resp)) }) } } -impl<'data, T> Drop for IoUringRuntime<'data, T> { +impl<'data, T, Meta> Drop for IoUringRuntime<'data, T, Meta> { fn drop(&mut self) { while self.in_progress > 0 || !self.io_uring.submission().is_empty() { // TODO: Cancel operations with `io_uring::Submitter::register_sync_cancel`? @@ -140,11 +140,11 @@ impl<'data, T> Drop for IoUringRuntime<'data, T> { } #[derive(Debug)] -pub struct IoUringState<'data, T> { - requests: Slab<(RequestId, IoUringRequest<'data, T>)>, +pub struct IoUringState<'data, T, Meta> { + requests: Slab<(Meta, IoUringRequest<'data, T>)>, } -impl<'data, T> IoUringState<'data, T> { +impl<'data, T, Meta> IoUringState<'data, T, Meta> { pub fn with_capacity(capacity: usize) -> Self { Self { requests: Slab::with_capacity(capacity), @@ -153,14 +153,7 @@ impl<'data, T> IoUringState<'data, T> { /// Allocates `Vec>`, reinterprets it as `Vec>`, and stores the byte buffer /// so the kernel writes into correctly aligned memory for `T`. - pub fn read( - &mut self, - id: RequestId, - fd: Fd, - range: ReadRange, - file_index: FileIndex, - direct_io: bool, - ) -> squeue::Entry + pub fn read(&mut self, meta: Meta, fd: Fd, range: ReadRange, direct_io: bool) -> squeue::Entry where T: bytemuck::Pod, { @@ -172,14 +165,7 @@ impl<'data, T> IoUringState<'data, T> { let mut items: Vec> = Vec::with_capacity(length as _); items.resize_with(length as _, || MaybeUninit::uninit()); - let (slot, req) = self.init( - id, - IoUringRequest::Read { - items, - file_index, - direct_io, - }, - ); + let (slot, req) = self.init(meta, IoUringRequest::Read { items, direct_io }); let items = req.expect_read(); let bytes_ptr = items.as_mut_ptr().cast(); @@ -193,7 +179,7 @@ impl<'data, T> IoUringState<'data, T> { pub fn write( &mut self, - id: RequestId, + meta: Meta, fd: Fd, byte_offset: u64, items: &'data [T], @@ -201,7 +187,7 @@ impl<'data, T> IoUringState<'data, T> { where T: bytemuck::Pod, { - let (slot, req) = self.init(id, IoUringRequest::Write(items)); + let (slot, req) = self.init(meta, IoUringRequest::Write(items)); let items = req.expect_write(); let bytes: &[u8] = bytemuck::cast_slice(items); @@ -214,12 +200,12 @@ impl<'data, T> IoUringState<'data, T> { fn init( &mut self, - id: RequestId, + meta: Meta, req: IoUringRequest<'data, T>, ) -> (usize, &mut IoUringRequest<'data, T>) { let entry = self.requests.vacant_entry(); let slot = entry.key(); - let (_, req) = entry.insert((id, req)); + let (_, req) = entry.insert((meta, req)); (slot, req) } @@ -227,8 +213,8 @@ impl<'data, T> IoUringState<'data, T> { &mut self, slot: usize, byte_length: u32, - ) -> io::Result<(RequestId, IoUringResponse)> { - let (id, req) = self + ) -> io::Result<(Meta, IoUringResponse)> { + let (meta, req) = self .requests .try_remove(slot) .ok_or_else(|| io::Error::other(format!("request in slot {slot} does not exist")))?; @@ -238,7 +224,6 @@ impl<'data, T> IoUringState<'data, T> { let resp = match req { IoUringRequest::Read { mut items, - file_index, direct_io, } => { if direct_io { @@ -251,7 +236,7 @@ impl<'data, T> IoUringState<'data, T> { } let items: Vec = unsafe { maybe_uninit::assume_init_vec(items) }; - IoUringResponse::Read { items, file_index } + IoUringResponse::Read(items) } IoUringRequest::Write(items) => { @@ -260,7 +245,7 @@ impl<'data, T> IoUringState<'data, T> { } }; - Ok((id, resp)) + Ok((meta, resp)) } pub fn abort(&mut self, slot: usize) { @@ -268,19 +253,16 @@ impl<'data, T> IoUringState<'data, T> { } } -impl<'data, T> Drop for IoUringState<'data, T> { +impl<'data, T, Meta> Drop for IoUringState<'data, T, Meta> { fn drop(&mut self) { debug_assert!(self.requests.is_empty()); } } -pub type RequestId = u64; - #[derive(Debug)] pub enum IoUringRequest<'data, T> { Read { items: Vec>, - file_index: FileIndex, direct_io: bool, }, @@ -307,19 +289,15 @@ impl<'data, T> IoUringRequest<'data, T> { #[derive(Debug)] pub enum IoUringResponse { - Read { - items: Vec, - file_index: FileIndex, - }, - + Read(Vec), Write, } impl IoUringResponse { - pub fn expect_read(self) -> (FileIndex, Vec) { + pub fn expect_read(self) -> Vec { #[expect(clippy::match_wildcard_for_single_variants)] match self { - Self::Read { items, file_index } => (file_index, items), + Self::Read(items) => items, _ => panic!(), } } diff --git a/lib/common/common/src/universal_io/io_uring/tests.rs b/lib/common/common/src/universal_io/io_uring/tests.rs index f4bac96dc0..9a6cbfd264 100644 --- a/lib/common/common/src/universal_io/io_uring/tests.rs +++ b/lib/common/common/src/universal_io/io_uring/tests.rs @@ -50,7 +50,7 @@ fn test_io_uring_read_batch() -> Result<()> { // Non-contiguous ranges across the file. #[rustfmt::skip] - let ranges = vec![ + let ranges = [ ReadRange { byte_offset: 0, length: 10 }, // [0..10] ReadRange { byte_offset: 50 * elem, length: 20 }, // [50..70] ReadRange { byte_offset: 100 * elem, length: 5 }, // [100..105] @@ -66,7 +66,7 @@ fn test_io_uring_read_batch() -> Result<()> { // --- read_batch (callback API) --- let mut batch_results: Vec<(usize, Vec)> = Vec::new(); - file.read_batch::(ranges.clone(), |idx, slice| { + file.read_batch::(ranges.iter().copied().enumerate(), |idx, slice| { batch_results.push((idx, slice.to_vec())); Ok(()) })?; @@ -82,7 +82,7 @@ fn test_io_uring_read_batch() -> Result<()> { // --- read_iter (iterator API) --- let mut iter_results: Vec<(usize, Vec)> = Vec::new(); - for record in file.read_iter::(ranges.clone()) { + for record in file.read_iter::(ranges.iter().copied().enumerate()) { let (idx, cow) = record?; iter_results.push((idx, cow.into_owned())); } @@ -97,15 +97,13 @@ fn test_io_uring_read_batch() -> Result<()> { } // --- read_iter with more ranges than the io_uring queue depth (64 > 16) --- - let many_ranges: Vec = (0..64) - .map(|i| ReadRange { - byte_offset: i * elem, - length: 1, - }) - .collect(); + let many_ranges = (0..64).map(|i| ReadRange { + byte_offset: i as u64 * elem, + length: 1, + }); let mut count = 0; - for record in file.read_iter::(many_ranges) { + for record in file.read_iter::(many_ranges.enumerate()) { let (idx, cow) = record?; assert_eq!( cow.as_ref(), @@ -148,16 +146,17 @@ fn test_io_uring_concurrent_read_iter() -> Result<()> { let file_b = TypedStorage::::open(&path_b, opts)?; // NUM_RANGES ranges, each reading CHUNK elements — well over the queue depth. - let ranges_a: Vec = (0..NUM_RANGES) - .map(|i| ReadRange { - byte_offset: i * CHUNK * elem, - length: CHUNK, - }) - .collect(); - let ranges_b: Vec = ranges_a.clone(); + let ranges_a = (0..NUM_RANGES).map(|i| ReadRange { + byte_offset: i * CHUNK * elem, + length: CHUNK, + }); + let ranges_b = (0..NUM_RANGES).map(|i| ReadRange { + byte_offset: i * CHUNK * elem, + length: CHUNK, + }); - let iter_a = file_a.read_iter::(ranges_a); - let iter_b = file_b.read_iter::(ranges_b); + let iter_a = file_a.read_iter::(ranges_a.enumerate()); + let iter_b = file_b.read_iter::(ranges_b.enumerate()); // Zip alternates next() calls between the two iterators on the same // thread-local io_uring ring. With in-flight operations left across @@ -209,31 +208,29 @@ fn test_io_uring_read_multi_iter_basic() -> Result<()> { // Interleaved reads across both files. #[rustfmt::skip] - let reads = vec![ - (0, ReadRange { byte_offset: 0, length: 10 }), // f0[0..10] - (1, ReadRange { byte_offset: 20 * elem, length: 5 }), // f1[20..25] - (0, ReadRange { byte_offset: 50 * elem, length: 20 }), // f0[50..70] - (1, ReadRange { byte_offset: 0, length: 10 }), // f1[0..10] + let reads = [ + ('a', &files[0], ReadRange { byte_offset: 0, length: 10 }), // f0[0..10] + ('b', &files[1], ReadRange { byte_offset: 20 * elem, length: 5 }), // f1[20..25] + ('c', &files[0], ReadRange { byte_offset: 50 * elem, length: 20 }), // f0[50..70] + ('d', &files[1], ReadRange { byte_offset: 0, length: 10 }), // f1[0..10] ]; - let expected: Vec<(FileIndex, &[u64])> = vec![ - (0, &data_0[0..10]), - (1, &data_1[20..25]), - (0, &data_0[50..70]), - (1, &data_1[0..10]), + let expected = [ + ('a', data_0[0..10].to_vec()), + ('b', data_1[20..25].to_vec()), + ('c', data_0[50..70].to_vec()), + ('d', data_1[0..10].to_vec()), ]; - let mut results: Vec<(usize, FileIndex, Vec)> = Vec::new(); - for record in IoUringFile::read_multi_iter::(&files, reads) { - let (idx, file_idx, cow) = record?; - results.push((idx, file_idx, cow.into_owned())); + let mut results: Vec<(char, Vec)> = Vec::new(); + for record in IoUringFile::read_multi_iter::(reads) { + let (idx, cow) = record?; + results.push((idx, cow.into_owned())); } - results.sort_by_key(|(idx, _, _)| *idx); - for (idx, file_idx, items) in &results { - let (expected_file, expected_data) = expected[*idx]; - assert_eq!(*file_idx, expected_file, "file index mismatch at op {idx}"); - assert_eq!(items.as_slice(), expected_data, "data mismatch at op {idx}"); + results.sort_by_key(|(idx, _)| *idx); + for (result, expected) in std::iter::zip(&results, &expected) { + assert_eq!(result, expected, "mismatch for read index {}", result.0); } Ok(()) @@ -263,13 +260,15 @@ fn test_io_uring_read_multi_iter_many_ranges() -> Result<()> { } // Generate reads: round-robin across files, each reading a small chunk. - let reads: Vec<(FileIndex, ReadRange)> = (0..NUM_FILES as u64 * RANGES_PER_FILE) + let reads: Vec<((usize, usize), &IoUringFile, ReadRange)> = (0..NUM_FILES as u64 + * RANGES_PER_FILE) .map(|i| { let file_idx = (i as usize) % NUM_FILES; let range_idx = i / NUM_FILES as u64; let offset = range_idx * 10; // non-overlapping chunks of 10 ( - file_idx, + (file_idx, offset as usize), + &files[file_idx], ReadRange { byte_offset: offset * elem, length: 10, @@ -278,24 +277,20 @@ fn test_io_uring_read_multi_iter_many_ranges() -> Result<()> { }) .collect(); - let mut results: Vec<(usize, FileIndex, Vec)> = Vec::new(); - for record in IoUringFile::read_multi_iter::(&files, reads.clone()) { - let (idx, file_idx, cow) = record?; - results.push((idx, file_idx, cow.into_owned())); + let mut results: Vec<((usize, usize), Vec)> = Vec::new(); + for record in IoUringFile::read_multi_iter::(reads) { + let (idx, cow) = record?; + results.push((idx, cow.into_owned())); } - assert_eq!(results.len(), reads.len()); + assert_eq!(results.len(), NUM_FILES * RANGES_PER_FILE as usize); - results.sort_by_key(|(idx, _, _)| *idx); - for (idx, file_idx, items) in &results { - let (expected_file, expected_range) = &reads[*idx]; - assert_eq!(file_idx, expected_file); - let start = (expected_range.byte_offset / elem) as usize; - let end = start + expected_range.length as usize; + results.sort_by_key(|(idx, _)| *idx); + for ((file_idx, offset), result) in results { assert_eq!( - items.as_slice(), - &all_data[*file_idx][start..end], - "data mismatch at op {idx}, file {file_idx}" + result.as_slice(), + &all_data[file_idx][offset..offset + 10], + "data mismatch at offset {offset}, file {file_idx}" ); } @@ -323,36 +318,35 @@ fn test_io_uring_read_multi_callback_matches_iter() -> Result<()> { let files = [file_a, file_b]; #[rustfmt::skip] - let reads: Vec<(FileIndex, ReadRange)> = vec![ - (0, ReadRange { byte_offset: 0, length: 50 }), - (1, ReadRange { byte_offset: 10 * elem, length: 30 }), - (0, ReadRange { byte_offset: 100 * elem, length: 50 }), - (1, ReadRange { byte_offset: 0, length: 100 }), - (0, ReadRange { byte_offset: 150 * elem, length: 50 }), + let reads: Vec<(usize, &IoUringFile, ReadRange)> = vec![ + (0, &files[0], ReadRange { byte_offset: 0, length: 50 }), + (1, &files[1], ReadRange { byte_offset: 10 * elem, length: 30 }), + (2, &files[0], ReadRange { byte_offset: 100 * elem, length: 50 }), + (3, &files[1], ReadRange { byte_offset: 0, length: 100 }), + (4, &files[0], ReadRange { byte_offset: 150 * elem, length: 50 }), ]; // Collect via callback. - let mut callback_results: Vec<(usize, FileIndex, Vec)> = Vec::new(); - IoUringFile::read_multi::(&files, reads.clone(), |idx, file_idx, data| { - callback_results.push((idx, file_idx, data.to_vec())); + let mut callback_results: Vec<(usize, Vec)> = Vec::new(); + IoUringFile::read_multi::(reads.clone(), |idx, data| { + callback_results.push((idx, data.to_vec())); Ok(()) })?; // Collect via iterator. - let mut iter_results: Vec<(usize, FileIndex, Vec)> = Vec::new(); - for record in IoUringFile::read_multi_iter::(&files, reads) { - let (idx, file_idx, cow) = record?; - iter_results.push((idx, file_idx, cow.into_owned())); + let mut iter_results: Vec<(usize, Vec)> = Vec::new(); + for record in IoUringFile::read_multi_iter::(reads) { + let (idx, cow) = record?; + iter_results.push((idx, cow.into_owned())); } - callback_results.sort_by_key(|(idx, _, _)| *idx); - iter_results.sort_by_key(|(idx, _, _)| *idx); + callback_results.sort_by_key(|(idx, _)| *idx); + iter_results.sort_by_key(|(idx, _)| *idx); assert_eq!(callback_results.len(), iter_results.len()); for (cb, it) in callback_results.iter().zip(iter_results.iter()) { assert_eq!(cb.0, it.0, "operation index mismatch"); - assert_eq!(cb.1, it.1, "file index mismatch"); - assert_eq!(cb.2, it.2, "data mismatch at op {}", cb.0); + assert_eq!(cb.1, it.1, "data mismatch at op {}", cb.0); } Ok(()) diff --git a/lib/common/common/src/universal_io/mmap.rs b/lib/common/common/src/universal_io/mmap.rs index 8714b6ef25..37cb433527 100644 --- a/lib/common/common/src/universal_io/mmap.rs +++ b/lib/common/common/src/universal_io/mmap.rs @@ -75,16 +75,16 @@ where Ok(Cow::Borrowed(items)) } - fn read_batch( - &self, - ranges: impl IntoIterator, - mut callback: impl FnMut(usize, &[T]) -> Result<()>, + fn read_batch<'a, P: AccessPattern, Meta: 'a>( + &'a self, + ranges: impl IntoIterator, + mut callback: impl FnMut(Meta, &[T]) -> Result<()>, ) -> Result<()> { let mmap = self.as_bytes::

(); - for (idx, range) in ranges.into_iter().enumerate() { + for (meta, range) in ranges { let items = read(mmap, range)?; - callback(idx, items)?; + callback(meta, items)?; } Ok(()) diff --git a/lib/common/common/src/universal_io/read.rs b/lib/common/common/src/universal_io/read.rs index 8d1ec27009..eae251bbc2 100644 --- a/lib/common/common/src/universal_io/read.rs +++ b/lib/common/common/src/universal_io/read.rs @@ -26,22 +26,21 @@ pub trait UniversalRead: UniversalReadFileOps { }) } - fn read_batch( - &self, - ranges: impl IntoIterator, - callback: impl FnMut(usize, &[T]) -> Result<()>, + fn read_batch<'a, P: AccessPattern, Meta: 'a>( + &'a self, + ranges: impl IntoIterator, + callback: impl FnMut(Meta, &[T]) -> Result<()>, ) -> Result<()>; /// Like [`read_batch`](Self::read_batch), but returns a fallible iterator instead of /// accepting a callback. - fn read_iter( + fn read_iter( &self, - ranges: impl IntoIterator, - ) -> impl Iterator)>> { + ranges: impl IntoIterator, + ) -> impl Iterator)>> { ranges .into_iter() - .enumerate() - .map(move |(idx, range)| self.read::

(range).map(|data| (idx, data))) + .map(move |(meta, range)| self.read::

(range).map(|data| (meta, data))) } fn len(&self) -> Result; @@ -57,21 +56,16 @@ pub trait UniversalRead: UniversalReadFileOps { fn clear_ram_cache(&self) -> Result<()>; /// Read from multiple files in a single operation. - fn read_multi( - files: &[Self], - reads: impl IntoIterator, - mut callback: impl FnMut(usize, FileIndex, &[T]) -> Result<()>, - ) -> Result<()> { - for (operation_index, (file_index, range)) in reads.into_iter().enumerate() { - let file = files - .get(file_index) - .ok_or(UniversalIoError::InvalidFileIndex { - file_index, - files: files.len(), - })?; - + fn read_multi<'a, P: AccessPattern, Meta: 'a>( + reads: impl IntoIterator, + mut callback: impl FnMut(Meta, &[T]) -> Result<()>, + ) -> Result<()> + where + Self: 'a, + { + for (meta, file, range) in reads { let data = file.read::

(range)?; - callback(operation_index, file_index, &data)?; + callback(meta, &data)?; } Ok(()) @@ -79,24 +73,16 @@ pub trait UniversalRead: UniversalReadFileOps { /// Like [`read_multi`](Self::read_multi), but returns a fallible iterator instead of /// accepting a callback. - fn read_multi_iter( - files: &[Self], - reads: impl IntoIterator, - ) -> impl Iterator)>> { - reads - .into_iter() - .enumerate() - .map(move |(idx, (file_index, range))| { - let file = files - .get(file_index) - .ok_or(UniversalIoError::InvalidFileIndex { - file_index, - files: files.len(), - })?; - - let data = file.read::

(range)?; - Ok((idx, file_index, data)) - }) + fn read_multi_iter<'a, P: AccessPattern, Meta>( + reads: impl IntoIterator, + ) -> impl Iterator)>> + where + Self: 'a, + { + reads.into_iter().map(move |(meta, file, range)| { + let data = file.read::

(range)?; + Ok((meta, data)) + }) } // When adding provided methods, don't forget to update impls in crate::universal_io::wrappers::*. diff --git a/lib/common/common/src/universal_io/wrappers/read_only.rs b/lib/common/common/src/universal_io/wrappers/read_only.rs index 9e41ce782e..d6697745e9 100644 --- a/lib/common/common/src/universal_io/wrappers/read_only.rs +++ b/lib/common/common/src/universal_io/wrappers/read_only.rs @@ -1,16 +1,10 @@ use std::borrow::Cow; use std::path::{Path, PathBuf}; -use bytemuck::TransparentWrapper; - -use super::super::{ - FileIndex, OpenOptions, ReadRange, Result, UniversalRead, UniversalReadFileOps, -}; +use super::super::{OpenOptions, ReadRange, Result, UniversalRead, UniversalReadFileOps}; use crate::generic_consts::AccessPattern; -#[derive(Debug, TransparentWrapper)] -#[repr(transparent)] -#[transparent(S)] +#[derive(Debug)] pub struct ReadOnly(S); impl UniversalReadFileOps for ReadOnly @@ -51,20 +45,20 @@ where } #[inline] - fn read_batch( - &self, - ranges: impl IntoIterator, - callback: impl FnMut(usize, &[T]) -> Result<()>, + fn read_batch<'a, P: AccessPattern, Meta: 'a>( + &'a self, + ranges: impl IntoIterator, + callback: impl FnMut(Meta, &[T]) -> Result<()>, ) -> Result<()> { - self.0.read_batch::

(ranges, callback) + self.0.read_batch::(ranges, callback) } #[inline] - fn read_iter( - &self, - ranges: impl IntoIterator, - ) -> impl Iterator)>> { - self.0.read_iter::

(ranges) + fn read_iter<'a, P: AccessPattern, Meta>( + &'a self, + ranges: impl IntoIterator, + ) -> impl Iterator)>> { + self.0.read_iter::(ranges) } #[inline] @@ -83,19 +77,29 @@ where } #[inline] - fn read_multi( - files: &[Self], - reads: impl IntoIterator, - callback: impl FnMut(usize, FileIndex, &[T]) -> Result<()>, - ) -> Result<()> { - S::read_multi::

(Self::peel_slice(files), reads, callback) + fn read_multi<'a, P: AccessPattern, Meta: 'a>( + reads: impl IntoIterator, + callback: impl FnMut(Meta, &[T]) -> Result<()>, + ) -> Result<()> + where + Self: 'a, + { + let reads = reads + .into_iter() + .map(|(meta, file, range)| (meta, &file.0, range)); + S::read_multi::(reads, callback) } #[inline] - fn read_multi_iter( - files: &[Self], - reads: impl IntoIterator, - ) -> impl Iterator)>> { - S::read_multi_iter::

(Self::peel_slice(files), reads) + fn read_multi_iter<'a, P: AccessPattern, Meta>( + reads: impl IntoIterator, + ) -> impl Iterator)>> + where + Self: 'a, + { + let it = reads + .into_iter() + .map(|(meta, file, range)| (meta, &file.0, range)); + S::read_multi_iter::(it) } } diff --git a/lib/common/common/src/universal_io/wrappers/typed.rs b/lib/common/common/src/universal_io/wrappers/typed.rs index 33194956a7..9a2e03f993 100644 --- a/lib/common/common/src/universal_io/wrappers/typed.rs +++ b/lib/common/common/src/universal_io/wrappers/typed.rs @@ -57,20 +57,20 @@ impl, T: Copy + 'static> UniversalRead for TypedStorage( - &self, - ranges: impl IntoIterator, - callback: impl FnMut(usize, &[T]) -> Result<()>, + fn read_batch<'a, P: AccessPattern, Meta: 'a>( + &'a self, + ranges: impl IntoIterator, + callback: impl FnMut(Meta, &[T]) -> Result<()>, ) -> Result<()> { - self.inner.read_batch::

(ranges, callback) + self.inner.read_batch::(ranges, callback) } #[inline] - fn read_iter( + fn read_iter( &self, - ranges: impl IntoIterator, - ) -> impl Iterator)>> { - self.inner.read_iter::

(ranges) + ranges: impl IntoIterator, + ) -> impl Iterator)>> { + self.inner.read_iter::(ranges) } #[inline] @@ -89,20 +89,32 @@ impl, T: Copy + 'static> UniversalRead for TypedStorage( - files: &[Self], - reads: impl IntoIterator, - callback: impl FnMut(usize, FileIndex, &[T]) -> Result<()>, - ) -> Result<()> { - S::read_multi::

(Self::peel_slice(files), reads, callback) + fn read_multi<'a, P: AccessPattern, Meta: 'a>( + reads: impl IntoIterator, + callback: impl FnMut(Meta, &[T]) -> Result<()>, + ) -> Result<()> + where + Self: 'a, + { + S::read_multi::<'a, P, Meta>( + reads + .into_iter() + .map(|(meta, file, range)| (meta, &file.inner, range)), + callback, + ) } #[inline] - fn read_multi_iter( - files: &[Self], - reads: impl IntoIterator, - ) -> impl Iterator)>> { - S::read_multi_iter::

(Self::peel_slice(files), reads) + fn read_multi_iter<'a, P: AccessPattern, Meta>( + reads: impl IntoIterator, + ) -> impl Iterator)>> + where + Self: 'a, + { + let reads = reads + .into_iter() + .map(|(meta, file, range)| (meta, &file.inner, range)); + S::read_multi_iter::(reads) } } diff --git a/lib/gridstore/src/pages.rs b/lib/gridstore/src/pages.rs index ccebdae6f4..f54e81461a 100644 --- a/lib/gridstore/src/pages.rs +++ b/lib/gridstore/src/pages.rs @@ -7,15 +7,11 @@ use common::maybe_uninit::assume_init_vec; use common::universal_io::{ FileIndex, Flusher, OpenOptions, ReadRange, UniversalRead, UniversalWrite, }; -use smallvec::SmallVec; use crate::Result; use crate::config::StorageConfig; use crate::tracker::{PageId, ValuePointer}; -type PageRanges = SmallVec<[(FileIndex, ReadRange); 2]>; -type BufferOffsets = SmallVec<[usize; 2]>; - pub fn page_path(base_path: &Path, page_id: PageId) -> PathBuf { base_path.join(format!("page_{page_id}.dat")) } @@ -80,24 +76,53 @@ impl> Pages { page_path(&self.base_path, page_id) } - /// Computes the page ranges and relative offsets for reading or writing a value described by - /// the given [`ValuePointer`]. + /// Computes the page ranges for reading or writing a value described by the + /// given [`ValuePointer`]. /// - /// A value may span across two consecutive pages if it starts near the end of a page. This - /// method returns the list of `(page_index, range)` pairs that together cover the full extent - /// of the value, along with the corresponding byte offsets into the value buffer so that each - /// section can be read into or written from the correct position. + /// Returns an iterator of tuples that together cover the full extent of the + /// value. Each tuple contains information about its chunk. /// - /// Returns a tuple of: - /// - `SmallVec<[(FileIndex, ReadRange); 2]>` — the per-page file ranges to access. - /// - `SmallVec<[usize; 2]>` — the offset within the value buffer that each page range - /// corresponds to. + /// Typically, the iterator contains only one entry (if value is within a + /// single page), or two entries if the value spans across two consecutive + /// pages. + /// + /// # Example + /// + /// Assume each page is 8×1K blocks. The requested [`ValuePointer`] spans + /// across two pages, so this function returns two tuples. + /// + /// ```text + /// ┌── ValuePointer ───┐ + /// │ page_id: 123 │ + /// │ block_offset: 6 │ ← requested value + /// │ length: 5K │ + /// └─────────┬─────────┘ + /// ╭──────────┴──────────╮ + /// ┐ ┌───┬───┬── page 123 ───┬───┬───┐ ┌───┬───┬── page 124 ───┬───┬───┐ ┌ + /// │ │ 0 │ 1 │ 2 │ 3 │ 4 │ 5 │ 6 │ 7 │ │ 0 │ 1 │ 2 │ 3 │ 4 │ 5 │ 6 │ 7 │ │ + /// ┘ └───┴───┴───┴───┴───┴───┴───┴───┘ └───┴───┴───┴───┴───┴───┴───┴───┘ └ + /// ╰───┬───╯ ╰─────┬─────╯ + /// ┌────── tuple 0 ─┴───┐ ┌─────┴ tuple 1 ─────┐ + /// │ buf_offset: 0 │ │ buf_offset: 2K │ ← returned + /// │ page_id: 123 │ │ page_id: 124 │ tuples + /// │ read_range: 6K..8K │ │ read_range: 0K..2K │ + /// └─────────────────┬──┘ └────┬───────────────┘ + /// ┌───┴───┬─────┴─────┐ + /// │ 0..2K │ 2K..5K │ ← value buffer + /// └───────┴───────────┘ + /// ``` fn get_page_value_ranges( pointer: ValuePointer, config: &StorageConfig, - ) -> (PageRanges, BufferOffsets) { + ) -> impl Iterator< + Item = ( + usize, // buf_offset - byte offset within the value buffer + PageId, + ReadRange, + ), + > { let ValuePointer { - page_id, + mut page_id, block_offset, length, } = pointer; @@ -109,35 +134,28 @@ impl> Pages { let value_start = u64::from(block_offset) * block_size_bytes; assert!(value_start < page_len); - // Do not expect payload to span more than 2 pages, but can be easily extended if needed - let mut page_ranges = PageRanges::new(); - // Store relative offsets to fill the raw_value buffer after fetching in batch - let mut buffer_offsets = BufferOffsets::new(); - - let mut length_so_far: u64 = 0; - let mut page_idx = page_id as FileIndex; + let mut buf_offset = 0; let mut start = value_start; - while length_so_far < total_length { - let remaining = total_length - length_so_far; + std::iter::from_fn(move || { + if buf_offset >= total_length { + return None; + } + let remaining = total_length - buf_offset; let available_on_page = page_len - start; let section_len = remaining.min(available_on_page); - page_ranges.push(( - page_idx, - ReadRange { - byte_offset: start, - length: section_len, - }, - )); - buffer_offsets.push(length_so_far as usize); + let read_range = ReadRange { + byte_offset: start, + length: section_len, + }; + let result = (buf_offset as usize, page_id, read_range); - length_so_far += section_len; - page_idx += 1; + buf_offset += section_len; + page_id += 1; start = 0; - } - - (page_ranges, buffer_offsets) + Some(result) + }) } pub fn read_from_pages( @@ -148,10 +166,10 @@ impl> Pages { // Avoid initializing buffer with zeros, as it will be overwritten by file access; let mut raw_value = vec![MaybeUninit::::uninit(); pointer.length as usize]; - let (read_ranges, buffer_offsets) = Self::get_page_value_ranges(pointer, config); + let reads = Self::get_page_value_ranges(pointer, config) + .map(|(buf_offset, page, range)| (buf_offset, &self.pages[page as usize], range)); - S::read_multi::

(self.pages.as_slice(), read_ranges, |idx, _, slice| { - let offset = buffer_offsets[idx]; + S::read_multi::(reads, |offset, slice| { raw_value[offset..offset + slice.len()].write_copy_of_slice(slice); Ok(()) })?; @@ -214,18 +232,10 @@ impl> Pages { value: &[u8], config: &StorageConfig, ) -> Result<()> { - // Compute write ranges - let (page_ranges, buffer_offsets) = Self::get_page_value_ranges(pointer, config); - - let writes = page_ranges - .iter() - .zip(&buffer_offsets) - .map(|(&(file_idx, range), &start)| { - ( - file_idx, - range.byte_offset, - &value[start..start + range.length as usize], - ) + let writes = + Self::get_page_value_ranges(pointer, config).map(|(buf_offset, page, range)| { + let data = &value[buf_offset..buf_offset + range.length as usize]; + (page as FileIndex, range.byte_offset, data) }); // Execute writes (mutable borrow of self.pages) diff --git a/lib/segment/src/vector_storage/dense/immutable_dense_vectors.rs b/lib/segment/src/vector_storage/dense/immutable_dense_vectors.rs index ddd7797ea3..b9f5bc0a25 100644 --- a/lib/segment/src/vector_storage/dense/immutable_dense_vectors.rs +++ b/lib/segment/src/vector_storage/dense/immutable_dense_vectors.rs @@ -197,11 +197,12 @@ impl> ImmutableDenseVectors length: self.dim as _, }); - self.storage.read_batch::

(ranges, |idx, vector| { - let point = points.get(idx).copied().expect("point ID tracked"); - callback(idx, point, vector); - Ok(()) - })?; + self.storage + .read_batch::(ranges.enumerate(), |idx, vector| { + let point = points.get(idx).copied().expect("point ID tracked"); + callback(idx, point, vector); + Ok(()) + })?; Ok(()) } diff --git a/src/tonic/api/storage_read_api/mod.rs b/src/tonic/api/storage_read_api/mod.rs index 4a7958fdf0..0b90fdeb65 100644 --- a/src/tonic/api/storage_read_api/mod.rs +++ b/src/tonic/api/storage_read_api/mod.rs @@ -296,7 +296,7 @@ impl + Send + Sync + 'static> StorageRead for StorageReadSe let storage = S::open(&path, open_options).map_err(io_error_to_status)?; let mut results = ranges.iter().map(|_| Vec::new()).collect::>(); storage - .read_batch::(ranges, |idx, chunk| { + .read_batch::(ranges.into_iter().enumerate(), |idx, chunk| { results[idx].extend_from_slice(chunk); Ok(()) }) @@ -356,7 +356,13 @@ impl + Send + Sync + 'static> StorageRead for StorageReadSe .map_err(io_error_to_status)?; let mut results = vec![Vec::new(); reads_.len()]; - S::read_multi::(&files, reads_, |op_idx, _, chunk| { + + let reads = reads_ + .into_iter() + .enumerate() + .map(|(op_idx, (file_idx, range))| (op_idx, &files[file_idx], range)); + + S::read_multi::(reads, |op_idx, chunk| { results[op_idx].extend_from_slice(chunk); Ok(()) })