tuwunel_database/
stream.rs1mod items;
2mod items_rev;
3mod keys;
4mod keys_rev;
5
6use std::{mem::replace, sync::Arc};
7
8use rocksdb::{DBRawIteratorWithThreadMode, ReadOptions};
9use tuwunel_core::Result;
10
11pub(crate) use self::{items::Items, items_rev::ItemsRev, keys::Keys, keys_rev::KeysRev};
12use crate::{
13 Map, Slice,
14 engine::Db,
15 keyval::{Key, KeyVal, Val},
16 util::{is_incomplete, map_err},
17};
18
19pub(crate) struct State<'a> {
25 inner: Inner<'a>,
26 seek: bool,
27 init: bool,
28}
29
30pub(crate) trait Cursor<'a, T>: Send {
36 fn state(&self) -> &State<'a>;
37
38 fn state_mut(&mut self) -> &mut State<'a>;
39
40 fn count(&self) -> (usize, Option<usize>);
41
42 fn fetch(&self) -> Option<T>;
43
44 fn seek(&mut self);
45
46 #[inline]
47 fn get(&self) -> Option<Result<T>> {
48 self.fetch()
49 .map(Ok)
50 .or_else(|| self.state().status().map(map_err).map(Err))
51 }
52
53 #[inline]
54 fn seek_and_get(&mut self) -> Option<Result<T>> {
55 self.seek();
56 self.get()
57 }
58}
59
60type Inner<'a> = DBRawIteratorWithThreadMode<'a, Db>;
61type From<'a> = Option<Key<'a>>;
62
63impl<'a> State<'a> {
64 #[inline]
65 pub(super) fn new(map: &'a Arc<Map>, opts: ReadOptions) -> Self {
66 Self {
67 init: true,
68 seek: false,
69 inner: map
70 .engine()
71 .db
72 .raw_iterator_cf_opt(&map.cf(), opts),
73 }
74 }
75
76 #[inline]
77 #[tracing::instrument(level = "trace", skip_all)]
78 pub(super) fn init_fwd(mut self, from: From<'_>) -> Self {
79 debug_assert!(self.init, "init must be set to make this call");
80 debug_assert!(!self.seek, "seek must not be set to make this call");
81
82 if let Some(key) = from {
83 self.inner.seek(key);
84 } else {
85 self.inner.seek_to_first();
86 }
87
88 self.seek = true;
89 self
90 }
91
92 #[inline]
93 #[tracing::instrument(level = "trace", skip_all)]
94 pub(super) fn init_rev(mut self, from: From<'_>) -> Self {
95 debug_assert!(self.init, "init must be set to make this call");
96 debug_assert!(!self.seek, "seek must not be set to make this call");
97
98 if let Some(key) = from {
99 self.inner.seek_for_prev(key);
100 } else {
101 self.inner.seek_to_last();
102 }
103
104 self.seek = true;
105 self
106 }
107
108 #[inline]
109 #[cfg_attr(unabridged, tracing::instrument(level = "trace", skip_all))]
110 pub(super) fn seek_fwd(&mut self) {
111 if !replace(&mut self.init, false) {
112 self.inner.next();
113 } else if !self.seek {
114 self.inner.seek_to_first();
115 }
116 }
117
118 #[inline]
119 #[cfg_attr(unabridged, tracing::instrument(level = "trace", skip_all))]
120 pub(super) fn seek_rev(&mut self) {
121 if !replace(&mut self.init, false) {
122 self.inner.prev();
123 } else if !self.seek {
124 self.inner.seek_to_last();
125 }
126 }
127
128 #[inline]
129 #[expect(clippy::unused_self)]
130 #[tracing::instrument(level = "trace", skip_all, ret)]
131 pub(super) fn count_fwd(&self) -> (usize, Option<usize>) { (0, None) }
132
133 #[inline]
134 #[expect(clippy::unused_self)]
135 #[tracing::instrument(level = "trace", skip_all, ret)]
136 pub(super) fn count_rev(&self) -> (usize, Option<usize>) { (0, None) }
137
138 #[inline]
139 fn fetch_key(&self) -> Option<Key<'_>> { self.inner.key() }
140
141 #[inline]
142 fn _fetch_val(&self) -> Option<Val<'_>> { self.inner.value() }
143
144 #[inline]
145 fn fetch(&self) -> Option<KeyVal<'_>> { self.inner.item() }
146
147 pub(super) fn is_incomplete(&self) -> bool {
148 matches!(self.status(), Some(e) if is_incomplete(&e))
149 }
150
151 #[inline]
152 pub(super) fn status(&self) -> Option<rocksdb::Error> { self.inner.status().err() }
153
154 #[inline]
155 pub(super) fn valid(&self) -> bool { self.inner.valid() }
156}
157
158fn keyval_longevity<'a, 'b: 'a>(item: KeyVal<'a>) -> KeyVal<'b> {
164 (slice_longevity::<'a, 'b>(item.0), slice_longevity::<'a, 'b>(item.1))
165}
166
167fn slice_longevity<'a, 'b: 'a>(item: &'a Slice) -> &'b Slice {
173 unsafe { std::mem::transmute(item) }
189}