use crate::{
error::{Error, Result},
query::{Query, QueryResult},
record::RecordRef,
utils::Path,
Library, Session,
};
use async_std::sync::Arc;
use std::{
collections::BTreeSet,
fmt::{self, Debug, Formatter},
mem,
sync::atomic::{AtomicUsize, Ordering},
};
pub struct QueryIterator {
pos: AtomicUsize,
paths: Vec<(Path, Session)>,
inner: Arc<Library>,
query: Query,
}
impl Debug for QueryIterator {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
write!(f, "{:?}", self.paths)
}
}
impl QueryIterator {
pub(crate) fn new(id: Session, paths: Vec<Path>, inner: Arc<Library>, query: Query) -> Self {
Self {
pos: 0.into(),
paths: paths.into_iter().map(|p| (p, id)).collect(),
inner,
query,
}
}
pub fn merge(&mut self, mut other: Self) -> Result<()> {
if self.query != other.query {
return Err(Error::IncompatibleQuery {
q1: format!("{:?}", self.query),
q2: format!("{:?}", other.query),
});
}
self.paths.append(&mut other.paths);
self.paths = mem::replace(&mut self.paths, vec![])
.into_iter()
.fold(BTreeSet::new(), |mut set, pp| {
set.insert(pp);
set
})
.into_iter()
.collect();
Ok(())
}
pub fn skip(&self, pos: usize) {
self.pos.fetch_add(pos, Ordering::Relaxed);
}
pub fn query(&self) -> &Query {
&self.query
}
#[inline]
pub fn pos(&self) -> usize {
self.pos.load(Ordering::Relaxed)
}
#[inline]
pub fn len(&self) -> usize {
self.paths.len()
}
#[inline]
pub fn remaining(&self) -> usize {
self.len() - self.pos()
}
pub async fn lock(&self) {
let mut s = self.inner.store.write().await;
s.gc_lock(&self.paths);
}
pub async fn next(&self) -> Result<Option<RecordRef>> {
if self.pos.load(Ordering::Relaxed) >= self.paths.len() {
return Ok(None);
}
let pos = self.pos.fetch_add(1, Ordering::Relaxed);
let (path, id) = self.paths.get(pos).unwrap().clone();
self.inner
.query(id, Query::Path(path))
.await
.map(|r| match r {
QueryResult::Single(rec) => Some(rec),
QueryResult::Many(_) => unreachable!(),
})
}
}
impl Drop for QueryIterator {
fn drop(&mut self) {
async_std::task::block_on(async {
let mut s = self.inner.store.write().await;
s.gc_release(&self.paths)
.expect("Failed to release deleted records!");
});
}
}
#[cfg(test)]
mod harness {
pub use crate::GLOBAL;
use crate::{
utils::{Diff, Path, TagSet},
Builder, Library,
};
use async_std::sync::Arc;
use hex;
use rand::{rngs::OsRng, RngCore};
use tempfile::tempdir;
pub struct TestData {
lib: Arc<Library>,
rng: OsRng,
}
impl TestData {
pub fn setup() -> Self {
let dir = tempdir().unwrap();
let lib = Builder::new().offset(dir.path()).build().unwrap();
let rng = OsRng {};
Self { lib, rng }
}
pub fn lib(&self) -> Arc<Library> {
Arc::clone(&self.lib)
}
pub async fn insert_random(&mut self) -> Path {
let mut seed = [0 as u8; 8];
self.rng.fill_bytes(&mut seed);
let name = hex::encode_upper(&seed);
let path = Path::from(format!("/test:{}", name));
self.rng.fill_bytes(&mut seed);
let key = hex::encode_upper(&seed);
self.rng.fill_bytes(&mut seed);
let value = hex::encode_upper(&seed);
self.lib
.insert(
GLOBAL,
path.clone(),
TagSet::empty(),
Diff::map().insert(key, value),
)
.await
.unwrap();
path
}
}
}
#[cfg(test)]
use harness::TestData;
#[async_std::test]
async fn basic_iterator() -> Result<()> {
let mut t = TestData::setup();
let paths = vec![
t.insert_random().await,
t.insert_random().await,
t.insert_random().await,
];
let iter = QueryIterator::new(harness::GLOBAL, paths, t.lib(), Query::Fake);
assert!(iter.next().await?.is_some());
assert!(iter.next().await?.is_some());
assert!(iter.next().await?.is_some());
assert!(iter.next().await?.is_none());
Ok(())
}
#[async_std::test]
async fn skip_iterator() -> Result<()> {
let mut t = TestData::setup();
let paths = vec![
t.insert_random().await,
t.insert_random().await,
t.insert_random().await,
];
let iter = QueryIterator::new(harness::GLOBAL, paths, t.lib(), Query::Fake);
iter.skip(2);
assert!(iter.next().await?.is_some());
assert!(iter.next().await?.is_none());
Ok(())
}
#[async_std::test]
async fn gc_iterator_fail() -> Result<()> {
let mut t = TestData::setup();
let paths = vec![
t.insert_random().await,
t.insert_random().await,
t.insert_random().await,
];
let iter = QueryIterator::new(harness::GLOBAL, paths.clone(), t.lib(), Query::Fake);
t.lib()
.delete(harness::GLOBAL, paths[0].clone())
.await
.unwrap();
assert!(iter.next().await.is_err());
Ok(())
}
#[async_std::test]
async fn lock_gc_iterator() -> Result<()> {
let mut t = TestData::setup();
let paths = vec![
t.insert_random().await,
t.insert_random().await,
t.insert_random().await,
];
let iter = QueryIterator::new(harness::GLOBAL, paths.clone(), t.lib(), Query::Fake);
iter.lock().await;
t.lib()
.delete(harness::GLOBAL, paths[0].clone())
.await
.unwrap();
assert!(iter.next().await?.is_some());
Ok(())
}