use std::ops::Range;
use super::api::{
BoxFut, KeySpan, KeyValSpan, KeyVisitor, KeysBatch, KeysResult, ScanChunkStats, ScanCursorKeys,
ScanCursorVals, ScanResult, Transactable, ValVisitor, ValsBatch,
};
use super::direction::Direction;
use super::err::Result;
use super::util;
use crate::kvs::{Key, Val};
async fn fetch_vals_page<T: Transactable + ?Sized>(
tx: &T,
rng: &mut Range<Key>,
dir: Direction,
version: Option<u64>,
skip: &mut u32,
exhausted: &mut bool,
limit: u32,
) -> Result<ScanResult> {
if *exhausted || rng.start >= rng.end {
return Ok(ScanResult::default());
}
let skip = std::mem::take(skip);
let res = match dir {
Direction::Forward => tx.scan(rng.clone(), limit, skip, version).await?,
Direction::Backward => tx.scanr(rng.clone(), limit, skip, version).await?,
};
match res.values.last() {
Some((last, _)) => {
util::update_range(rng, dir, Some(last));
if res.values.len() < limit as usize {
*exhausted = true;
}
}
None => *exhausted = true,
}
Ok(res)
}
#[allow(clippy::too_many_arguments, reason = "threads the cursor's scan state + reusable arenas")]
pub(crate) async fn fill_vals_batch<T: Transactable + ?Sized>(
tx: &T,
rng: &mut Range<Key>,
dir: Direction,
version: Option<u64>,
skip: &mut u32,
exhausted: &mut bool,
key_buf: &mut Vec<u8>,
val_buf: &mut Vec<u8>,
spans: &mut Vec<KeyValSpan>,
limit: u32,
) -> Result<(u64, u64)> {
key_buf.clear();
val_buf.clear();
spans.clear();
let res = fetch_vals_page(tx, rng, dir, version, skip, exhausted, limit).await?;
let (kb, vb): (usize, usize) =
res.values.iter().fold((0, 0), |(ka, va), (k, v)| (ka + k.len(), va + v.len()));
key_buf.reserve(kb);
val_buf.reserve(vb);
spans.reserve(res.values.len());
for (k, v) in &res.values {
let key_offset = key_buf.len();
let key_len = k.len();
key_buf.extend_from_slice(k);
let val_offset = val_buf.len();
let val_len = v.len();
val_buf.extend_from_slice(v);
spans.push(KeyValSpan {
key_offset,
key_len,
val_offset,
val_len,
});
}
Ok((res.key_bytes, res.value_bytes))
}
async fn fetch_keys_page<T: Transactable + ?Sized>(
tx: &T,
rng: &mut Range<Key>,
dir: Direction,
version: Option<u64>,
skip: &mut u32,
exhausted: &mut bool,
limit: u32,
) -> Result<KeysResult> {
if *exhausted || rng.start >= rng.end {
return Ok(KeysResult::default());
}
let skip = std::mem::take(skip);
let res = match dir {
Direction::Forward => tx.keys(rng.clone(), limit, skip, version).await?,
Direction::Backward => tx.keysr(rng.clone(), limit, skip, version).await?,
};
match res.keys.last() {
Some(last) => {
util::update_range(rng, dir, Some(last));
if res.keys.len() < limit as usize {
*exhausted = true;
}
}
None => *exhausted = true,
}
Ok(res)
}
#[allow(clippy::too_many_arguments, reason = "threads the cursor's scan state + reusable arena")]
async fn fill_keys_batch<T: Transactable + ?Sized>(
tx: &T,
rng: &mut Range<Key>,
dir: Direction,
version: Option<u64>,
skip: &mut u32,
exhausted: &mut bool,
key_buf: &mut Vec<u8>,
key_spans: &mut Vec<KeySpan>,
limit: u32,
) -> Result<u64> {
key_buf.clear();
key_spans.clear();
let res = fetch_keys_page(tx, rng, dir, version, skip, exhausted, limit).await?;
let total_bytes: usize = res.keys.iter().map(|k| k.len()).sum();
key_buf.reserve(total_bytes);
key_spans.reserve(res.keys.len());
for k in &res.keys {
let offset = key_buf.len();
let len = k.len();
key_buf.extend_from_slice(k);
key_spans.push(KeySpan {
offset,
len,
});
}
Ok(res.key_bytes)
}
pub(crate) struct DefaultKeysCursor<'a, T: ?Sized> {
tx: &'a T,
rng: Range<Key>,
dir: Direction,
version: Option<u64>,
skip: u32,
exhausted: bool,
key_buf: Vec<u8>,
key_spans: Vec<KeySpan>,
pending: std::vec::IntoIter<Key>,
}
impl<'a, T: ?Sized> DefaultKeysCursor<'a, T> {
pub(crate) fn new(
tx: &'a T,
rng: Range<Key>,
dir: Direction,
version: Option<u64>,
skip: u32,
) -> Self {
Self {
tx,
rng,
dir,
version,
skip,
exhausted: false,
key_buf: Vec::new(),
key_spans: Vec::new(),
pending: Vec::new().into_iter(),
}
}
}
impl<T> ScanCursorKeys for DefaultKeysCursor<'_, T>
where
T: Transactable + ?Sized,
{
fn next_batch<'s>(&'s mut self, limit: u32) -> BoxFut<'s, Result<KeysBatch<'s>>> {
Box::pin(async move {
debug_assert!(
self.pending.as_slice().is_empty(),
"next_batch called while for_each left rows buffered (early Break or mid-page limit stop); the two paths must not be mixed on one cursor",
);
let key_bytes = fill_keys_batch(
self.tx,
&mut self.rng,
self.dir,
self.version,
&mut self.skip,
&mut self.exhausted,
&mut self.key_buf,
&mut self.key_spans,
limit,
)
.await?;
Ok(KeysBatch::from_parts(&self.key_buf, &self.key_spans, key_bytes))
})
}
fn for_each<'s>(
&'s mut self,
limit: u32,
f: &'s mut dyn KeyVisitor,
) -> BoxFut<'s, Result<ScanChunkStats>> {
Box::pin(async move {
let mut stats = ScanChunkStats::default();
loop {
while stats.rows < limit as u64 {
let Some(k) = self.pending.next() else {
break;
};
let flow = f(&k)?;
stats.rows += 1;
stats.key_bytes += k.len() as u64;
if let std::ops::ControlFlow::Break(()) = flow {
return Ok(stats);
}
}
if stats.rows >= limit as u64 {
return Ok(stats);
}
let res = fetch_keys_page(
self.tx,
&mut self.rng,
self.dir,
self.version,
&mut self.skip,
&mut self.exhausted,
limit,
)
.await?;
if res.keys.is_empty() {
return Ok(stats);
}
self.pending = res.keys.into_iter();
}
})
}
}
pub(crate) struct DefaultValsCursor<'a, T: ?Sized> {
tx: &'a T,
rng: Range<Key>,
dir: Direction,
version: Option<u64>,
skip: u32,
exhausted: bool,
key_buf: Vec<u8>,
val_buf: Vec<u8>,
spans: Vec<KeyValSpan>,
pending: std::vec::IntoIter<(Key, Val)>,
}
impl<'a, T: ?Sized> DefaultValsCursor<'a, T> {
pub(crate) fn new(
tx: &'a T,
rng: Range<Key>,
dir: Direction,
version: Option<u64>,
skip: u32,
) -> Self {
Self {
tx,
rng,
dir,
version,
skip,
exhausted: false,
key_buf: Vec::new(),
val_buf: Vec::new(),
spans: Vec::new(),
pending: Vec::new().into_iter(),
}
}
}
impl<T> ScanCursorVals for DefaultValsCursor<'_, T>
where
T: Transactable + ?Sized,
{
fn next_batch<'s>(&'s mut self, limit: u32) -> BoxFut<'s, Result<ValsBatch<'s>>> {
Box::pin(async move {
debug_assert!(
self.pending.as_slice().is_empty(),
"next_batch called while for_each left rows buffered (early Break or mid-page limit stop); the two paths must not be mixed on one cursor",
);
let (key_bytes, value_bytes) = fill_vals_batch(
self.tx,
&mut self.rng,
self.dir,
self.version,
&mut self.skip,
&mut self.exhausted,
&mut self.key_buf,
&mut self.val_buf,
&mut self.spans,
limit,
)
.await?;
Ok(ValsBatch::from_parts(
&self.key_buf,
&self.val_buf,
&self.spans,
key_bytes,
value_bytes,
))
})
}
fn for_each<'s>(
&'s mut self,
limit: u32,
f: &'s mut dyn ValVisitor,
) -> BoxFut<'s, Result<ScanChunkStats>> {
Box::pin(async move {
let mut stats = ScanChunkStats::default();
loop {
while stats.rows < limit as u64 {
let Some((k, v)) = self.pending.next() else {
break;
};
let flow = f(&k, &v)?;
stats.rows += 1;
stats.key_bytes += k.len() as u64;
stats.value_bytes += v.len() as u64;
if let std::ops::ControlFlow::Break(()) = flow {
return Ok(stats);
}
}
if stats.rows >= limit as u64 {
return Ok(stats);
}
let res = fetch_vals_page(
self.tx,
&mut self.rng,
self.dir,
self.version,
&mut self.skip,
&mut self.exhausted,
limit,
)
.await?;
if res.values.is_empty() {
return Ok(stats);
}
self.pending = res.values.into_iter();
}
})
}
}