use std::marker::PhantomPinned;
use std::ops::RangeBounds;
use std::pin::Pin;
use std::sync::Arc;
use std::task::Context;
use std::task::Poll;
use futures::Stream;
use crate::v001::block::Block;
use crate::v001::block::BlockIter;
use crate::v001::SegmentedKey;
use crate::v001::SeqMarked;
pub struct BlockStream {
iter: BlockIter<'static>,
#[allow(dead_code)]
block: Arc<Block>,
_p: PhantomPinned,
}
impl BlockStream {
pub fn new<R>(block: Arc<Block>, range: R) -> Self
where R: RangeBounds<String> {
let block_ptr = block.as_ref() as *const Block;
let block_ref = unsafe { &*block_ptr };
let iter = block_ref.range(range);
Self {
block,
iter,
_p: Default::default(),
}
}
fn next(self: Pin<&mut Self>) -> Option<(SegmentedKey, &SeqMarked)> {
let it = unsafe { &mut self.get_unchecked_mut().iter };
it.next()
}
}
impl Stream for BlockStream {
type Item = (String, SeqMarked);
fn poll_next(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
let next = self.next().map(|(k, v)| (k.to_string(), v.clone()));
Poll::Ready(next)
}
}
#[cfg(test)]
#[allow(clippy::redundant_clone)]
mod tests {
use std::sync::Arc;
use futures::executor::block_on;
use futures::StreamExt;
use crate::v001::block::Block;
use crate::v001::block_stream::BlockStream;
use crate::v001::testing::bb;
use crate::v001::testing::ss;
use crate::v001::testing::ss_vec;
use crate::v001::SeqMarked;
#[test]
fn test_block_stream() -> anyhow::Result<()> {
let block_data = maplit::btreemap! {
ss("a") => SeqMarked::new_tombstone(1),
ss("b") => SeqMarked::new_normal(2, bb("B")),
ss("c") => SeqMarked::new_normal(3, bb("C")),
ss("d") => SeqMarked::new_normal(4, bb("D")),
};
let block = Block::new(5, block_data.clone());
let block = Arc::new(block);
{
let stream = BlockStream::new(block.clone(), ..);
let block_ptr = stream.block.as_ref() as *const Block;
println!("block_ptr: {:x}", block_ptr as usize);
}
fn collect(strm: BlockStream) -> Vec<String> {
block_on(strm.map(|(k, _v)| k).collect::<Vec<_>>())
}
{
let stream = BlockStream::new(block.clone(), ..);
let got = collect(stream);
assert_eq!(ss_vec(["a", "b", "c", "d"]), got);
}
{
let stream = BlockStream::new(block.clone(), ..ss("a"));
let got = collect(stream);
assert_eq!(Vec::<String>::new(), got);
}
{
let stream = BlockStream::new(block.clone(), ss("b1")..);
let got = collect(stream);
assert_eq!(ss_vec(["c", "d"]), got);
}
{
let stream = BlockStream::new(block.clone(), ..ss("c1"));
let got = collect(stream);
assert_eq!(ss_vec(["a", "b", "c"]), got);
}
{
let stream = BlockStream::new(block.clone(), ss("b1")..ss("c1"));
let got = collect(stream);
assert_eq!(ss_vec(["c"]), got);
}
Ok(())
}
}