rocksdb/db_iterator.rs
1// Copyright 2020 Tyler Neely
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7// http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use crate::{
16 db::{DBAccess, DB},
17 ffi, Error, ReadOptions, WriteBatch,
18};
19use libc::{c_char, c_uchar, size_t};
20use std::{marker::PhantomData, slice};
21
22/// A type alias to keep compatibility. See [`DBRawIteratorWithThreadMode`] for details
23pub type DBRawIterator<'a> = DBRawIteratorWithThreadMode<'a, DB>;
24
25/// A low-level iterator over a database or column family, created by [`DB::raw_iterator`]
26/// and other `raw_iterator_*` methods.
27///
28/// This iterator replicates RocksDB's API. It should provide better
29/// performance and more features than [`DBIteratorWithThreadMode`], which is a standard
30/// Rust [`std::iter::Iterator`].
31///
32/// ```
33/// use rocksdb::{DB, Options};
34///
35/// let tempdir = tempfile::Builder::new()
36/// .prefix("_path_for_rocksdb_storage4")
37/// .tempdir()
38/// .expect("Failed to create temporary path for the _path_for_rocksdb_storage4.");
39/// let path = tempdir.path();
40/// {
41/// let db = DB::open_default(path).unwrap();
42/// let mut iter = db.raw_iterator();
43///
44/// // Forwards iteration
45/// iter.seek_to_first();
46/// while iter.valid() {
47/// println!("Saw {:?} {:?}", iter.key(), iter.value());
48/// iter.next();
49/// }
50///
51/// // Reverse iteration
52/// iter.seek_to_last();
53/// while iter.valid() {
54/// println!("Saw {:?} {:?}", iter.key(), iter.value());
55/// iter.prev();
56/// }
57///
58/// // Seeking
59/// iter.seek(b"my key");
60/// while iter.valid() {
61/// println!("Saw {:?} {:?}", iter.key(), iter.value());
62/// iter.next();
63/// }
64///
65/// // Reverse iteration from key
66/// // Note, use seek_for_prev when reversing because if this key doesn't exist,
67/// // this will make the iterator start from the previous key rather than the next.
68/// iter.seek_for_prev(b"my key");
69/// while iter.valid() {
70/// println!("Saw {:?} {:?}", iter.key(), iter.value());
71/// iter.prev();
72/// }
73/// }
74/// let _ = DB::destroy(&Options::default(), path);
75/// ```
76pub struct DBRawIteratorWithThreadMode<'a, D: DBAccess> {
77 inner: std::ptr::NonNull<ffi::rocksdb_iterator_t>,
78
79 /// When iterate_lower_bound or iterate_upper_bound are set, the inner
80 /// C iterator keeps a pointer to the upper bound inside `_readopts`.
81 /// Storing this makes sure the upper bound is always alive when the
82 /// iterator is being used.
83 ///
84 /// And yes, we need to store the entire ReadOptions structure since C++
85 /// ReadOptions keep reference to C rocksdb_readoptions_t wrapper which
86 /// point to vectors we own. See issue #660.
87 _readopts: ReadOptions,
88
89 db: PhantomData<&'a D>,
90}
91
92impl<'a, D: DBAccess> DBRawIteratorWithThreadMode<'a, D> {
93 pub(crate) fn new(db: &D, readopts: ReadOptions) -> Self {
94 let inner = unsafe { db.create_iterator(&readopts) };
95 Self::from_inner(inner, readopts)
96 }
97
98 pub(crate) fn new_cf(
99 db: &'a D,
100 cf_handle: *mut ffi::rocksdb_column_family_handle_t,
101 readopts: ReadOptions,
102 ) -> Self {
103 let inner = unsafe { db.create_iterator_cf(cf_handle, &readopts) };
104 Self::from_inner(inner, readopts)
105 }
106
107 fn from_inner(inner: *mut ffi::rocksdb_iterator_t, readopts: ReadOptions) -> Self {
108 // This unwrap will never fail since rocksdb_create_iterator and
109 // rocksdb_create_iterator_cf functions always return non-null. They
110 // use new and dereference the result so any nulls would end up with SIGSEGV
111 // there and we would have a bigger issue.
112 let inner = std::ptr::NonNull::new(inner).unwrap();
113 Self {
114 inner,
115 _readopts: readopts,
116 db: PhantomData,
117 }
118 }
119
120 /// Returns `true` if the iterator is valid. An iterator is invalidated when
121 /// it reaches the end of its defined range, or when it encounters an error.
122 ///
123 /// To check whether the iterator encountered an error after `valid` has
124 /// returned `false`, use the [`status`](DBRawIteratorWithThreadMode::status) method. `status` will never
125 /// return an error when `valid` is `true`.
126 pub fn valid(&self) -> bool {
127 unsafe { ffi::rocksdb_iter_valid(self.inner.as_ptr()) != 0 }
128 }
129
130 /// Returns an error `Result` if the iterator has encountered an error
131 /// during operation. When an error is encountered, the iterator is
132 /// invalidated and [`valid`](DBRawIteratorWithThreadMode::valid) will return `false` when called.
133 ///
134 /// Performing a seek will discard the current status.
135 pub fn status(&self) -> Result<(), Error> {
136 unsafe {
137 ffi_try!(ffi::rocksdb_iter_get_error(self.inner.as_ptr()));
138 }
139 Ok(())
140 }
141
142 /// Refreshes the iterator to represent the latest state of the DB.
143 /// The iterator is invalidated after this call and must be re-sought
144 /// before use.
145 ///
146 /// If the iterator was created with a snapshot, the refreshed iterator
147 /// will no longer use that snapshot and will instead read the latest
148 /// DB state. The snapshot itself is not released; it remains valid and
149 /// will be released when the owning [`crate::SnapshotWithThreadMode`] is dropped.
150 pub fn refresh(&mut self) -> Result<(), Error> {
151 unsafe {
152 ffi_try!(ffi::rocksdb_iter_refresh(self.inner.as_ptr()));
153 }
154 Ok(())
155 }
156
157 /// Seeks to the first key in the database.
158 ///
159 /// # Examples
160 ///
161 /// ```rust
162 /// use rocksdb::{DB, Options};
163 ///
164 /// let tempdir = tempfile::Builder::new()
165 /// .prefix("_path_for_rocksdb_storage5")
166 /// .tempdir()
167 /// .expect("Failed to create temporary path for the _path_for_rocksdb_storage5.");
168 /// let path = tempdir.path();
169 /// {
170 /// let db = DB::open_default(path).unwrap();
171 /// let mut iter = db.raw_iterator();
172 ///
173 /// // Iterate all keys from the start in lexicographic order
174 /// iter.seek_to_first();
175 ///
176 /// while iter.valid() {
177 /// println!("{:?} {:?}", iter.key(), iter.value());
178 /// iter.next();
179 /// }
180 ///
181 /// // Read just the first key
182 /// iter.seek_to_first();
183 ///
184 /// if iter.valid() {
185 /// println!("{:?} {:?}", iter.key(), iter.value());
186 /// } else {
187 /// // There are no keys in the database
188 /// }
189 /// }
190 /// let _ = DB::destroy(&Options::default(), path);
191 /// ```
192 pub fn seek_to_first(&mut self) {
193 unsafe {
194 ffi::rocksdb_iter_seek_to_first(self.inner.as_ptr());
195 }
196 }
197
198 /// Seeks to the last key in the database.
199 ///
200 /// # Examples
201 ///
202 /// ```rust
203 /// use rocksdb::{DB, Options};
204 ///
205 /// let tempdir = tempfile::Builder::new()
206 /// .prefix("_path_for_rocksdb_storage6")
207 /// .tempdir()
208 /// .expect("Failed to create temporary path for the _path_for_rocksdb_storage6.");
209 /// let path = tempdir.path();
210 /// {
211 /// let db = DB::open_default(path).unwrap();
212 /// let mut iter = db.raw_iterator();
213 ///
214 /// // Iterate all keys from the end in reverse lexicographic order
215 /// iter.seek_to_last();
216 ///
217 /// while iter.valid() {
218 /// println!("{:?} {:?}", iter.key(), iter.value());
219 /// iter.prev();
220 /// }
221 ///
222 /// // Read just the last key
223 /// iter.seek_to_last();
224 ///
225 /// if iter.valid() {
226 /// println!("{:?} {:?}", iter.key(), iter.value());
227 /// } else {
228 /// // There are no keys in the database
229 /// }
230 /// }
231 /// let _ = DB::destroy(&Options::default(), path);
232 /// ```
233 pub fn seek_to_last(&mut self) {
234 unsafe {
235 ffi::rocksdb_iter_seek_to_last(self.inner.as_ptr());
236 }
237 }
238
239 /// Seeks to the specified key or the first key that lexicographically follows it.
240 ///
241 /// This method will attempt to seek to the specified key. If that key does not exist, it will
242 /// find and seek to the key that lexicographically follows it instead.
243 ///
244 /// # Examples
245 ///
246 /// ```rust
247 /// use rocksdb::{DB, Options};
248 ///
249 /// let tempdir = tempfile::Builder::new()
250 /// .prefix("_path_for_rocksdb_storage7")
251 /// .tempdir()
252 /// .expect("Failed to create temporary path for the _path_for_rocksdb_storage7.");
253 /// let path = tempdir.path();
254 /// {
255 /// let db = DB::open_default(path).unwrap();
256 /// let mut iter = db.raw_iterator();
257 ///
258 /// // Read the first key that starts with 'a'
259 /// iter.seek(b"a");
260 ///
261 /// if iter.valid() {
262 /// println!("{:?} {:?}", iter.key(), iter.value());
263 /// } else {
264 /// // There are no keys in the database
265 /// }
266 /// }
267 /// let _ = DB::destroy(&Options::default(), path);
268 /// ```
269 pub fn seek<K: AsRef<[u8]>>(&mut self, key: K) {
270 let key = key.as_ref();
271
272 unsafe {
273 ffi::rocksdb_iter_seek(
274 self.inner.as_ptr(),
275 key.as_ptr() as *const c_char,
276 key.len() as size_t,
277 );
278 }
279 }
280
281 /// Seeks to the specified key, or the first key that lexicographically precedes it.
282 ///
283 /// Like ``.seek()`` this method will attempt to seek to the specified key.
284 /// The difference with ``.seek()`` is that if the specified key do not exist, this method will
285 /// seek to key that lexicographically precedes it instead.
286 ///
287 /// # Examples
288 ///
289 /// ```rust
290 /// use rocksdb::{DB, Options};
291 ///
292 /// let tempdir = tempfile::Builder::new()
293 /// .prefix("_path_for_rocksdb_storage8")
294 /// .tempdir()
295 /// .expect("Failed to create temporary path for the _path_for_rocksdb_storage8.");
296 /// let path = tempdir.path();
297 /// {
298 /// let db = DB::open_default(path).unwrap();
299 /// let mut iter = db.raw_iterator();
300 ///
301 /// // Read the last key that starts with 'a'
302 /// iter.seek_for_prev(b"b");
303 ///
304 /// if iter.valid() {
305 /// println!("{:?} {:?}", iter.key(), iter.value());
306 /// } else {
307 /// // There are no keys in the database
308 /// }
309 /// }
310 /// let _ = DB::destroy(&Options::default(), path);
311 /// ```
312 pub fn seek_for_prev<K: AsRef<[u8]>>(&mut self, key: K) {
313 let key = key.as_ref();
314
315 unsafe {
316 ffi::rocksdb_iter_seek_for_prev(
317 self.inner.as_ptr(),
318 key.as_ptr() as *const c_char,
319 key.len() as size_t,
320 );
321 }
322 }
323
324 /// Seeks to the next key.
325 pub fn next(&mut self) {
326 if self.valid() {
327 unsafe {
328 ffi::rocksdb_iter_next(self.inner.as_ptr());
329 }
330 }
331 }
332
333 /// Seeks to the previous key.
334 pub fn prev(&mut self) {
335 if self.valid() {
336 unsafe {
337 ffi::rocksdb_iter_prev(self.inner.as_ptr());
338 }
339 }
340 }
341
342 /// Returns a slice of the current key.
343 pub fn key(&self) -> Option<&[u8]> {
344 if self.valid() {
345 Some(self.key_impl())
346 } else {
347 None
348 }
349 }
350
351 /// Returns a slice of the current value.
352 pub fn value(&self) -> Option<&[u8]> {
353 if self.valid() {
354 Some(self.value_impl())
355 } else {
356 None
357 }
358 }
359
360 /// Returns a slice of the timestamp of the current entry.
361 pub fn timestamp(&self) -> Option<&[u8]> {
362 if self.valid() {
363 Some(self.timestamp_impl())
364 } else {
365 None
366 }
367 }
368
369 /// Returns pair with slice of the current key and current value.
370 pub fn item(&self) -> Option<(&[u8], &[u8])> {
371 if self.valid() {
372 Some((self.key_impl(), self.value_impl()))
373 } else {
374 None
375 }
376 }
377
378 /// Returns a slice of the current key; assumes the iterator is valid.
379 fn key_impl(&self) -> &[u8] {
380 // Safety Note: This is safe as all methods that may invalidate the buffer returned
381 // take `&mut self`, so borrow checker will prevent use of buffer after seek.
382 unsafe {
383 let mut key_len: size_t = 0;
384 let key_len_ptr: *mut size_t = &raw mut key_len;
385 let key_ptr = ffi::rocksdb_iter_key(self.inner.as_ptr(), key_len_ptr);
386 slice::from_raw_parts(key_ptr as *const c_uchar, key_len)
387 }
388 }
389
390 /// Returns a slice of the current value; assumes the iterator is valid.
391 fn value_impl(&self) -> &[u8] {
392 // Safety Note: This is safe as all methods that may invalidate the buffer returned
393 // take `&mut self`, so borrow checker will prevent use of buffer after seek.
394 unsafe {
395 let mut val_len: size_t = 0;
396 let val_len_ptr: *mut size_t = &raw mut val_len;
397 let val_ptr = ffi::rocksdb_iter_value(self.inner.as_ptr(), val_len_ptr);
398 slice::from_raw_parts(val_ptr as *const c_uchar, val_len)
399 }
400 }
401
402 /// Returns a slice of the timestamp of the entry; assumes the iterator is valid.
403 fn timestamp_impl(&self) -> &[u8] {
404 // Safety Note: This is safe as all methods that may invalidate the buffer returned
405 // take `&mut self`, so borrow checker will prevent use of buffer after seek.
406 unsafe {
407 let mut timestamp_len: size_t = 0;
408 let timestamp_len_ptr: *mut size_t = &raw mut timestamp_len;
409 let timestamp_ptr = ffi::rocksdb_iter_timestamp(self.inner.as_ptr(), timestamp_len_ptr);
410 slice::from_raw_parts(timestamp_ptr as *const c_uchar, timestamp_len)
411 }
412 }
413}
414
415impl<D: DBAccess> Drop for DBRawIteratorWithThreadMode<'_, D> {
416 fn drop(&mut self) {
417 unsafe {
418 ffi::rocksdb_iter_destroy(self.inner.as_ptr());
419 }
420 }
421}
422
423unsafe impl<D: DBAccess> Send for DBRawIteratorWithThreadMode<'_, D> {}
424unsafe impl<D: DBAccess> Sync for DBRawIteratorWithThreadMode<'_, D> {}
425
426/// A type alias to keep compatibility. See [`DBIteratorWithThreadMode`] for details
427pub type DBIterator<'a> = DBIteratorWithThreadMode<'a, DB>;
428
429/// A standard Rust [`Iterator`] over a database or column family.
430///
431/// As an alternative, [`DBRawIteratorWithThreadMode`] is a low level wrapper around
432/// RocksDB's API, which can provide better performance and more features.
433///
434/// ```
435/// use rocksdb::{DB, Direction, IteratorMode, Options};
436///
437/// let tempdir = tempfile::Builder::new()
438/// .prefix("_path_for_rocksdb_storage2")
439/// .tempdir()
440/// .expect("Failed to create temporary path for the _path_for_rocksdb_storage2.");
441/// let path = tempdir.path();
442/// {
443/// let db = DB::open_default(path).unwrap();
444/// let mut iter = db.iterator(IteratorMode::Start); // Always iterates forward
445/// for item in iter {
446/// let (key, value) = item.unwrap();
447/// println!("Saw {:?} {:?}", key, value);
448/// }
449/// iter = db.iterator(IteratorMode::End); // Always iterates backward
450/// for item in iter {
451/// let (key, value) = item.unwrap();
452/// println!("Saw {:?} {:?}", key, value);
453/// }
454/// iter = db.iterator(IteratorMode::From(b"my key", Direction::Forward)); // From a key in Direction::{forward,reverse}
455/// for item in iter {
456/// let (key, value) = item.unwrap();
457/// println!("Saw {:?} {:?}", key, value);
458/// }
459///
460/// // You can seek with an existing Iterator instance, too
461/// iter = db.iterator(IteratorMode::Start);
462/// iter.set_mode(IteratorMode::From(b"another key", Direction::Reverse));
463/// for item in iter {
464/// let (key, value) = item.unwrap();
465/// println!("Saw {:?} {:?}", key, value);
466/// }
467/// }
468/// let _ = DB::destroy(&Options::default(), path);
469/// ```
470pub struct DBIteratorWithThreadMode<'a, D: DBAccess> {
471 raw: DBRawIteratorWithThreadMode<'a, D>,
472 direction: Direction,
473 done: bool,
474}
475
476#[derive(Copy, Clone)]
477pub enum Direction {
478 Forward,
479 Reverse,
480}
481
482pub type KVBytes = (Box<[u8]>, Box<[u8]>);
483
484#[derive(Copy, Clone)]
485pub enum IteratorMode<'a> {
486 Start,
487 End,
488 From(&'a [u8], Direction),
489}
490
491impl<'a, D: DBAccess> DBIteratorWithThreadMode<'a, D> {
492 pub(crate) fn new(db: &D, readopts: ReadOptions, mode: IteratorMode) -> Self {
493 Self::from_raw(DBRawIteratorWithThreadMode::new(db, readopts), mode)
494 }
495
496 pub(crate) fn new_cf(
497 db: &'a D,
498 cf_handle: *mut ffi::rocksdb_column_family_handle_t,
499 readopts: ReadOptions,
500 mode: IteratorMode,
501 ) -> Self {
502 Self::from_raw(
503 DBRawIteratorWithThreadMode::new_cf(db, cf_handle, readopts),
504 mode,
505 )
506 }
507
508 fn from_raw(raw: DBRawIteratorWithThreadMode<'a, D>, mode: IteratorMode) -> Self {
509 let mut rv = DBIteratorWithThreadMode {
510 raw,
511 direction: Direction::Forward, // blown away by set_mode()
512 done: false,
513 };
514 rv.set_mode(mode);
515 rv
516 }
517
518 pub fn set_mode(&mut self, mode: IteratorMode) {
519 self.done = false;
520 self.direction = match mode {
521 IteratorMode::Start => {
522 self.raw.seek_to_first();
523 Direction::Forward
524 }
525 IteratorMode::End => {
526 self.raw.seek_to_last();
527 Direction::Reverse
528 }
529 IteratorMode::From(key, Direction::Forward) => {
530 self.raw.seek(key);
531 Direction::Forward
532 }
533 IteratorMode::From(key, Direction::Reverse) => {
534 self.raw.seek_for_prev(key);
535 Direction::Reverse
536 }
537 };
538 }
539
540 /// Refreshes the iterator, then re-seeks using the given mode.
541 ///
542 /// After a refresh the underlying iterator is invalidated, so a mode
543 /// must be provided to reposition it.
544 pub fn refresh(&mut self, mode: IteratorMode) -> Result<(), Error> {
545 self.raw.refresh()?;
546 self.set_mode(mode);
547 Ok(())
548 }
549}
550
551impl<D: DBAccess> Iterator for DBIteratorWithThreadMode<'_, D> {
552 type Item = Result<KVBytes, Error>;
553
554 fn next(&mut self) -> Option<Result<KVBytes, Error>> {
555 if self.done {
556 None
557 } else if let Some((key, value)) = self.raw.item() {
558 let item = (Box::from(key), Box::from(value));
559 match self.direction {
560 Direction::Forward => self.raw.next(),
561 Direction::Reverse => self.raw.prev(),
562 }
563 Some(Ok(item))
564 } else {
565 self.done = true;
566 self.raw.status().err().map(Result::Err)
567 }
568 }
569}
570
571impl<D: DBAccess> std::iter::FusedIterator for DBIteratorWithThreadMode<'_, D> {}
572
573impl<'a, D: DBAccess> Into<DBRawIteratorWithThreadMode<'a, D>> for DBIteratorWithThreadMode<'a, D> {
574 fn into(self) -> DBRawIteratorWithThreadMode<'a, D> {
575 self.raw
576 }
577}
578
579/// Iterates the batches of writes since a given sequence number.
580///
581/// `DBWALIterator` is returned by `DB::get_updates_since()` and will return the
582/// batches of write operations that have occurred since a given sequence number
583/// (see `DB::latest_sequence_number()`). This iterator cannot be constructed by
584/// the application.
585///
586/// The iterator item type is a tuple of (`u64`, `WriteBatch`) where the first
587/// value is the sequence number of the associated write batch.
588///
589pub struct DBWALIterator {
590 pub(crate) inner: *mut ffi::rocksdb_wal_iterator_t,
591 pub(crate) start_seq_number: u64,
592}
593
594impl DBWALIterator {
595 /// Returns `true` if the iterator is valid. An iterator is invalidated when
596 /// it reaches the end of its defined range, or when it encounters an error.
597 ///
598 /// To check whether the iterator encountered an error after `valid` has
599 /// returned `false`, use the [`status`](DBWALIterator::status) method.
600 /// `status` will never return an error when `valid` is `true`.
601 pub fn valid(&self) -> bool {
602 unsafe { ffi::rocksdb_wal_iter_valid(self.inner) != 0 }
603 }
604
605 /// Returns an error `Result` if the iterator has encountered an error
606 /// during operation. When an error is encountered, the iterator is
607 /// invalidated and [`valid`](DBWALIterator::valid) will return `false` when
608 /// called.
609 pub fn status(&self) -> Result<(), Error> {
610 unsafe {
611 ffi_try!(ffi::rocksdb_wal_iter_status(self.inner));
612 }
613 Ok(())
614 }
615}
616
617impl Iterator for DBWALIterator {
618 type Item = Result<(u64, WriteBatch), Error>;
619
620 fn next(&mut self) -> Option<Self::Item> {
621 if !self.valid() {
622 return None;
623 }
624
625 let mut seq: u64 = 0;
626 let mut batch = WriteBatch {
627 inner: unsafe { ffi::rocksdb_wal_iter_get_batch(self.inner, &raw mut seq) },
628 };
629
630 // if the initial sequence number is what was requested we skip it to
631 // only provide changes *after* it
632 while seq <= self.start_seq_number {
633 unsafe {
634 ffi::rocksdb_wal_iter_next(self.inner);
635 }
636
637 if !self.valid() {
638 return None;
639 }
640
641 // this drops which in turn frees the skipped batch
642 batch = WriteBatch {
643 inner: unsafe { ffi::rocksdb_wal_iter_get_batch(self.inner, &raw mut seq) },
644 };
645 }
646
647 if !self.valid() {
648 return self.status().err().map(Result::Err);
649 }
650
651 // Seek to the next write batch.
652 // Note that WriteBatches live independently of the WAL iterator so this is safe to do
653 unsafe {
654 ffi::rocksdb_wal_iter_next(self.inner);
655 }
656
657 Some(Ok((seq, batch)))
658 }
659}
660
661impl Drop for DBWALIterator {
662 fn drop(&mut self) {
663 unsafe {
664 ffi::rocksdb_wal_iter_destroy(self.inner);
665 }
666 }
667}