Skip to main content

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}