use super::error::Error;
use super::node::ReadNode;
use super::read::EntryRef;
use super::store::Store;
use super::tree::Tree;
use super::Prolly;
use super::{store::AsyncStore, AsyncProlly};
use futures_util::stream::{self, Stream};
use serde::{Deserialize, Serialize};
use std::sync::Arc;
type RangeItem = Result<(Vec<u8>, Vec<u8>), Error>;
type LeafEntry = (Vec<u8>, Vec<u8>);
pub(crate) const RANGE_CHILD_PREFETCH_PARALLELISM: usize = 16;
#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct RangeCursor {
after_key: Option<Vec<u8>>,
}
impl RangeCursor {
pub fn start() -> Self {
Self { after_key: None }
}
pub fn after_key(key: impl Into<Vec<u8>>) -> Self {
Self {
after_key: Some(key.into()),
}
}
pub fn after(&self) -> Option<&[u8]> {
self.after_key.as_deref()
}
pub fn is_start(&self) -> bool {
self.after_key.is_none()
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct RangePage {
pub entries: Vec<(Vec<u8>, Vec<u8>)>,
pub next_cursor: Option<RangeCursor>,
}
#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ReverseCursor {
before_key: Option<Vec<u8>>,
}
impl ReverseCursor {
pub fn end() -> Self {
Self { before_key: None }
}
pub fn before_key(key: impl Into<Vec<u8>>) -> Self {
Self {
before_key: Some(key.into()),
}
}
pub fn before(&self) -> Option<&[u8]> {
self.before_key.as_deref()
}
pub fn is_end(&self) -> bool {
self.before_key.is_none()
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct ReversePage {
pub entries: Vec<(Vec<u8>, Vec<u8>)>,
pub next_cursor: Option<ReverseCursor>,
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct CursorWindow {
pub position_key: Option<Vec<u8>>,
pub position_value: Option<Vec<u8>>,
pub found: bool,
pub entries: Vec<(Vec<u8>, Vec<u8>)>,
pub next_cursor: Option<RangeCursor>,
}
pub type AsyncRangePage = RangePage;
pub type AsyncReversePage = ReversePage;
pub fn create_range_iter<'a, S: Store>(
prolly: &'a Prolly<S>,
tree: &Tree,
start: &[u8],
end: Option<&[u8]>,
) -> Result<RangeIter<'a, S>, Error> {
let ready_store = prolly.engine.store.clone();
let future = create_async_range_iter(&prolly.engine, tree, start, end);
let inner = super::engine::ready::run_ready(ready_store.ready(future))?;
Ok(RangeIter { inner })
}
pub fn create_range_after_iter<'a, S: Store>(
prolly: &'a Prolly<S>,
tree: &Tree,
after_key: &[u8],
end: Option<&[u8]>,
) -> Result<RangeIter<'a, S>, Error> {
let ready_store = prolly.engine.store.clone();
let future = create_async_range_after_iter(&prolly.engine, tree, after_key, end);
let inner = super::engine::ready::run_ready(ready_store.ready(future))?;
Ok(RangeIter { inner })
}
pub async fn create_async_range_iter<'a, S>(
prolly: &'a AsyncProlly<S>,
tree: &Tree,
start: &[u8],
end: Option<&[u8]>,
) -> Result<AsyncRangeIter<'a, S>, Error>
where
S: AsyncStore,
S::Error: Send + Sync,
{
if end.is_some_and(|end| end <= start) {
return Ok(AsyncRangeIter::new(prolly, Vec::new(), start, end));
}
let path = prolly.find_read_path_arcs(tree, start).await?;
Ok(AsyncRangeIter::new(prolly, path, start, end))
}
pub async fn create_async_range_after_iter<'a, S>(
prolly: &'a AsyncProlly<S>,
tree: &Tree,
after_key: &[u8],
end: Option<&[u8]>,
) -> Result<AsyncRangeIter<'a, S>, Error>
where
S: AsyncStore,
S::Error: Send + Sync,
{
if end.is_some_and(|end| end <= after_key) {
return Ok(AsyncRangeIter::new_after(
prolly,
Vec::new(),
after_key,
end,
));
}
let path = prolly.find_read_path_arcs(tree, after_key).await?;
Ok(AsyncRangeIter::new_after(prolly, path, after_key, end))
}
pub struct RangeIter<'a, S: Store> {
inner: AsyncRangeIter<'a, super::store::SyncStoreAsAsync<Arc<S>>>,
}
impl<S: Store> RangeIter<'_, S> {
pub fn resume_cursor(&self) -> RangeCursor {
self.inner.resume_cursor()
}
}
impl<S: Store> Iterator for RangeIter<'_, S> {
type Item = RangeItem;
fn next(&mut self) -> Option<Self::Item> {
if let (Some((node, index)), Some(leaf_end)) =
(self.inner.stack.last_mut(), self.inner.leaf_end)
{
if node.is_leaf() && *index < leaf_end {
let entry_index = *index;
*index += 1;
if let Some((_, last_index)) = self.inner.last_location.as_mut() {
*last_index = entry_index;
}
return Some(node.entry_owned(entry_index).ok_or(Error::InvalidNode));
}
}
loop {
match self.inner.ready_step() {
RangeStep::Entry => {
let Some((node, index)) = self.inner.last_location.as_ref() else {
return Some(Err(Error::InvalidNode));
};
return Some(node.entry_owned(*index).ok_or(Error::InvalidNode));
}
RangeStep::NeedsChild => {
let ready_store = self.inner.prolly.store.clone();
let future = self.inner.descend_next_child();
if let Err(error) = super::engine::ready::run_ready(ready_store.ready(future)) {
return Some(Err(error));
}
}
RangeStep::Finished => return None,
RangeStep::Error(error) => return Some(Err(error)),
}
}
}
}
enum RangeStep {
Entry,
NeedsChild,
Finished,
Error(Error),
}
pub struct AsyncRangeIter<'a, S: AsyncStore> {
prolly: &'a AsyncProlly<S>,
stack: Vec<(Arc<ReadNode>, usize)>,
end: Option<Vec<u8>>,
started: bool,
start_key: Vec<u8>,
skip_start_key: bool,
last_location: Option<(Arc<ReadNode>, usize)>,
leaf_end: Option<usize>,
}
impl<'a, S> AsyncRangeIter<'a, S>
where
S: AsyncStore,
S::Error: Send + Sync,
{
pub(crate) fn new(
prolly: &'a AsyncProlly<S>,
stack: Vec<(Arc<ReadNode>, usize)>,
start: &[u8],
end: Option<&[u8]>,
) -> Self {
Self {
prolly,
stack,
end: end.map(|e| e.to_vec()),
started: false,
start_key: start.to_vec(),
skip_start_key: false,
last_location: None,
leaf_end: None,
}
}
pub(crate) fn new_after(
prolly: &'a AsyncProlly<S>,
stack: Vec<(Arc<ReadNode>, usize)>,
after_key: &[u8],
end: Option<&[u8]>,
) -> Self {
Self {
prolly,
stack,
end: end.map(|e| e.to_vec()),
started: false,
start_key: after_key.to_vec(),
skip_start_key: true,
last_location: None,
leaf_end: None,
}
}
pub async fn next(&mut self) -> Option<RangeItem> {
if let (Some((node, index)), Some(leaf_end)) = (self.stack.last_mut(), self.leaf_end) {
if node.is_leaf() && *index < leaf_end {
let entry_index = *index;
*index += 1;
if let Some((_, last_index)) = self.last_location.as_mut() {
*last_index = entry_index;
}
return Some(node.entry_owned(entry_index).ok_or(Error::InvalidNode));
}
}
loop {
match self.ready_step() {
RangeStep::Entry => {
let Some((node, index)) = self.last_location.as_ref() else {
return Some(Err(Error::InvalidNode));
};
return Some(node.entry_owned(*index).ok_or(Error::InvalidNode));
}
RangeStep::NeedsChild => {
if let Err(error) = self.descend_next_child().await {
return Some(Err(error));
}
}
RangeStep::Finished => return None,
RangeStep::Error(error) => return Some(Err(error)),
}
}
}
pub async fn next_with<R>(
&mut self,
read: impl for<'entry> FnOnce(EntryRef<'entry>) -> R,
) -> Option<Result<R, Error>> {
let mut read = Some(read);
if let (Some((node, index)), Some(leaf_end)) = (self.stack.last_mut(), self.leaf_end) {
if node.is_leaf() && *index < leaf_end {
let entry_index = *index;
*index += 1;
if let Some((_, last_index)) = self.last_location.as_mut() {
*last_index = entry_index;
}
let Some((key, value)) = node.entry(entry_index) else {
return Some(Err(Error::InvalidNode));
};
return Some(Ok(read
.take()
.expect("async range callback is invoked at most once")(
EntryRef::new(key, value),
)));
}
}
loop {
match self.ready_step() {
RangeStep::Entry => {
let Some((node, index)) = self.last_location.as_ref() else {
return Some(Err(Error::InvalidNode));
};
let Some((key, value)) = node.entry(*index) else {
return Some(Err(Error::InvalidNode));
};
return Some(Ok(read
.take()
.expect("async range callback is invoked at most once")(
EntryRef::new(key, value),
)));
}
RangeStep::NeedsChild => {
if let Err(error) = self.descend_next_child().await {
return Some(Err(error));
}
}
RangeStep::Finished => return None,
RangeStep::Error(error) => return Some(Err(error)),
}
}
}
pub async fn collect(mut self) -> Result<Vec<LeafEntry>, Error> {
let mut entries = Vec::new();
while let Some(item) = self.next().await {
entries.push(item?);
}
Ok(entries)
}
pub fn resume_cursor(&self) -> RangeCursor {
self.last_location
.as_ref()
.and_then(|(node, index)| node.key(*index))
.map(<[u8]>::to_vec)
.map(RangeCursor::after_key)
.unwrap_or_else(RangeCursor::start)
}
pub fn into_stream(self) -> impl Stream<Item = RangeItem> + 'a {
stream::unfold(self, |mut iter| async move {
iter.next().await.map(|item| (item, iter))
})
}
fn position_at_start(&mut self) {
if self.started {
return;
}
self.started = true;
let Some((node, idx)) = self.stack.last_mut() else {
return;
};
if node.is_leaf() {
*idx = match node.search(&self.start_key) {
Ok(i) if self.skip_start_key => i.saturating_add(1),
Ok(i) | Err(i) => i,
};
self.leaf_end = Some(self.end.as_deref().map_or(node.len(), |end| {
node.search(end).unwrap_or_else(|index| index)
}));
}
}
fn ready_step(&mut self) -> RangeStep {
self.position_at_start();
loop {
let Some((node, idx)) = self.stack.last_mut() else {
return RangeStep::Finished;
};
if *idx >= node.len() {
match self.advance_to_next_sibling() {
Ok(true) => continue,
Ok(false) => return RangeStep::Finished,
Err(error) => return RangeStep::Error(error),
}
}
if node.is_leaf() {
let entering_leaf = self.leaf_end.is_none();
let leaf_end = match self.leaf_end {
Some(leaf_end) => leaf_end,
None => {
let leaf_end = self.end.as_deref().map_or(node.len(), |end| {
node.search(end).unwrap_or_else(|index| index)
});
self.leaf_end = Some(leaf_end);
leaf_end
}
};
if *idx >= leaf_end {
return RangeStep::Finished;
}
let index = *idx;
*idx += 1;
match self.last_location.as_mut() {
Some((_, last_index)) if !entering_leaf => *last_index = index,
_ => self.last_location = Some((node.clone(), index)),
}
return RangeStep::Entry;
}
return match child_starts_at_or_after_end(self.end.as_deref(), node, *idx) {
Ok(true) => RangeStep::Finished,
Ok(false) => RangeStep::NeedsChild,
Err(error) => RangeStep::Error(error),
};
}
}
async fn descend_next_child(&mut self) -> Result<(), Error> {
let (node, index) = self
.stack
.last()
.map(|(node, index)| (node.clone(), *index))
.ok_or(Error::InvalidNode)?;
let child = self.load_child_for_descent(&node, index).await?;
self.stack.push((child, 0));
self.leaf_end = None;
Ok(())
}
fn advance_to_next_sibling(&mut self) -> Result<bool, Error> {
loop {
self.stack.pop();
self.leaf_end = None;
let Some((parent, parent_idx)) = self.stack.last_mut() else {
return Ok(false);
};
*parent_idx += 1;
if *parent_idx < parent.len() {
if child_starts_at_or_after_end(self.end.as_deref(), parent, *parent_idx)? {
return Ok(false);
}
return Ok(true);
}
}
}
async fn load_child_for_descent(
&self,
node: &ReadNode,
idx: usize,
) -> Result<Arc<ReadNode>, Error> {
let child_cid = node.child_cid(idx)?;
if !self.prolly.store().prefers_batch_reads() {
return self.prolly.load_read_arc(&child_cid).await;
}
let max_child_idx = node
.len()
.min(idx.saturating_add(RANGE_CHILD_PREFETCH_PARALLELISM));
let mut child_cids = Vec::with_capacity(max_child_idx.saturating_sub(idx));
child_cids.push(child_cid);
for child_idx in idx + 1..max_child_idx {
if child_starts_at_or_after_end(self.end.as_deref(), node, child_idx).unwrap_or(true) {
break;
}
match node.child_cid(child_idx) {
Ok(cid) => child_cids.push(cid),
Err(_) => break,
}
}
if child_cids.len() == 1 {
return self.prolly.load_read_arc(&child_cids[0]).await;
}
let children = self.prolly.load_many_read_ordered(&child_cids).await?;
children.into_iter().next().ok_or(Error::InvalidNode)
}
}
fn child_starts_at_or_after_end(
end: Option<&[u8]>,
node: &ReadNode,
child_index: usize,
) -> Result<bool, Error> {
let Some(end) = end else {
return Ok(false);
};
let first_key = node.key(child_index).ok_or(Error::InvalidNode)?;
Ok(first_key >= end)
}