use super::tx::Transaction;
use super::Key;
use super::Val;
use crate::err::Error;
use crate::idx::planner::ScanDirection;
use futures::stream::Stream;
use futures::Future;
use futures::FutureExt;
use std::collections::VecDeque;
use std::ops::Range;
use std::pin::Pin;
use std::task::{Context, Poll};
#[cfg(not(target_family = "wasm"))]
type FutureResult<'a, I> = Pin<Box<dyn Future<Output = Result<Vec<I>, Error>> + 'a + Send>>;
#[cfg(target_family = "wasm")]
type FutureResult<'a, I> = Pin<Box<dyn Future<Output = Result<Vec<I>, Error>> + 'a>>;
pub(super) struct Scanner<'a, I> {
store: &'a Transaction,
batch: u32,
range: Range<Key>,
results: VecDeque<I>,
future: Option<FutureResult<'a, I>>,
exhausted: bool,
version: Option<u64>,
limit: Option<usize>,
sc: ScanDirection,
}
impl<'a, I> Scanner<'a, I> {
pub fn new(
store: &'a Transaction,
batch: u32,
range: Range<Key>,
version: Option<u64>,
limit: Option<usize>,
sc: ScanDirection,
) -> Self {
Scanner {
store,
batch,
range,
future: None,
results: VecDeque::new(),
exhausted: false,
version,
limit,
sc,
}
}
fn next_poll<S, K>(
&mut self,
cx: &mut Context,
scan: S,
key: K,
) -> Poll<Option<Result<I, Error>>>
where
S: Fn(Range<Key>, u32) -> FutureResult<'a, I>,
K: Fn(&I) -> &Key,
{
if let Some(v) = self.results.pop_front() {
return Poll::Ready(Some(Ok(v)));
}
if self.exhausted {
return Poll::Ready(None);
}
if self.future.is_none() {
let range = self.range.clone();
let batch = self
.limit
.map(|l| (self.batch as usize).min(l) as u32)
.unwrap_or_else(|| self.batch);
self.future = Some(scan(range, batch));
}
match self.future.as_mut().unwrap().poll_unpin(cx) {
Poll::Ready(result) => {
self.future = None;
match result {
Ok(v) => match v.is_empty() {
true => {
Poll::Ready(None)
}
false => {
if let Some(l) = &mut self.limit {
*l -= v.len();
}
if v.len() < self.batch as usize {
self.exhausted = true;
}
let last = v.last().ok_or_else(|| {
fail!("Expected the last key-value pair to not be none")
})?;
match self.sc {
ScanDirection::Forward => {
self.range.start.clone_from(key(last));
self.range.start.push(0xff);
}
#[cfg(any(feature = "kv-rocksdb", feature = "kv-tikv"))]
ScanDirection::Backward => {
self.range.end.clone_from(key(last));
}
};
self.results.extend(v);
let item = self.results.pop_front().unwrap();
Poll::Ready(Some(Ok(item)))
}
},
Err(error) => Poll::Ready(Some(Err(error))),
}
}
Poll::Pending => Poll::Pending,
}
}
}
impl Stream for Scanner<'_, (Key, Val)> {
type Item = Result<(Key, Val), Error>;
fn poll_next(
mut self: Pin<&mut Self>,
cx: &mut Context,
) -> Poll<Option<Result<(Key, Val), Error>>> {
let (store, version) = (self.store, self.version);
match self.sc {
ScanDirection::Forward => self.next_poll(
cx,
move |range, batch| Box::pin(store.scan(range, batch, version)),
|v| &v.0,
),
#[cfg(any(feature = "kv-rocksdb", feature = "kv-tikv"))]
ScanDirection::Backward => self.next_poll(
cx,
move |range, batch| Box::pin(store.scanr(range, batch, version)),
|v| &v.0,
),
}
}
}
impl Stream for Scanner<'_, Key> {
type Item = Result<Key, Error>;
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context) -> Poll<Option<Result<Key, Error>>> {
let (store, version) = (self.store, self.version);
match self.sc {
ScanDirection::Forward => self.next_poll(
cx,
move |range, batch| Box::pin(store.keys(range, batch, version)),
|v| v,
),
#[cfg(any(feature = "kv-rocksdb", feature = "kv-tikv"))]
ScanDirection::Backward => self.next_poll(
cx,
move |range, batch| Box::pin(store.keysr(range, batch, version)),
|v| v,
),
}
}
}