tuwunel_database/map/
seek.rs1use std::sync::Arc;
2
3use futures::{FutureExt, Stream, StreamExt, TryFutureExt, TryStreamExt};
4use rocksdb::Direction;
5use tokio::task::consume_budget;
6use tuwunel_core::Result;
7
8use super::{Map, cache_iter_options_default, iter_options_default};
9use crate::{
10 pool::{Seek, into_send_seek},
11 stream,
12};
13
14pub(super) fn seek_stream<'a, C, T>(
21 map: &'a Arc<Map>,
22 dir: Direction,
23 from: Option<&[u8]>,
24) -> impl Stream<Item = Result<T>> + Send + use<'a, C, T>
25where
26 C: From<stream::State<'a>> + Stream<Item = Result<T>> + Send,
27{
28 let opts = iter_options_default(&map.engine);
29 let state = stream::State::new(map, opts);
30 if is_cached(map, dir, from) {
31 let state = init(state, dir, from);
32 return consume_budget()
33 .map(move |()| C::from(state))
34 .into_stream()
35 .flatten()
36 .left_stream();
37 }
38
39 let seek = Seek {
40 map: map.clone(),
41 state: into_send_seek(state),
42 dir,
43 key: from.map(Into::into),
44 res: None,
45 };
46
47 map.engine
48 .pool
49 .execute_iter(seek)
50 .ok_into::<C>()
51 .into_stream()
52 .try_flatten()
53 .right_stream()
54}
55
56#[tracing::instrument(
61 name = "cached",
62 level = "trace",
63 skip_all,
64 fields(%map),
65)]
66fn is_cached(map: &Arc<Map>, dir: Direction, from: Option<&[u8]>) -> bool {
67 let opts = cache_iter_options_default(&map.engine);
68 let state = init(stream::State::new(map, opts), dir, from);
69
70 !state.is_incomplete()
71}
72
73fn init<'a>(state: stream::State<'a>, dir: Direction, from: Option<&[u8]>) -> stream::State<'a> {
78 match dir {
79 | Direction::Forward => state.init_fwd(from),
80 | Direction::Reverse => state.init_rev(from),
81 }
82}