Skip to main content

rocksdb/transactions/
transaction_db.rs

1// Copyright 2021 Yiyuan Liu
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//
15
16use std::{
17    collections::BTreeMap,
18    ffi::CString,
19    fs, iter,
20    marker::PhantomData,
21    path::{Path, PathBuf},
22    ptr,
23    sync::{Arc, Mutex},
24};
25
26use crate::CStrLike;
27use std::ffi::CStr;
28
29use crate::column_family::ColumnFamilyTtl;
30use crate::{
31    column_family::UnboundColumnFamily,
32    db::{convert_values, DBAccess},
33    db_options::OptionsMustOutliveDB,
34    ffi,
35    ffi_util::to_cpath,
36    AsColumnFamilyRef, BoundColumnFamily, ColumnFamily, ColumnFamilyDescriptor,
37    DBIteratorWithThreadMode, DBPinnableSlice, DBRawIteratorWithThreadMode, Direction, Error,
38    FlushOptions, IteratorMode, MultiThreaded, Options, ReadOptions, SingleThreaded,
39    SnapshotWithThreadMode, ThreadMode, Transaction, TransactionDBOptions, TransactionOptions,
40    WriteBatchWithTransaction, WriteOptions, DB, DEFAULT_COLUMN_FAMILY_NAME,
41};
42use ffi::rocksdb_transaction_t;
43use libc::{c_char, c_int, c_uchar, c_void, size_t};
44
45#[cfg(not(feature = "multi-threaded-cf"))]
46type DefaultThreadMode = crate::SingleThreaded;
47#[cfg(feature = "multi-threaded-cf")]
48type DefaultThreadMode = crate::MultiThreaded;
49
50/// RocksDB TransactionDB.
51///
52/// Please read the official [guide](https://github.com/facebook/rocksdb/wiki/Transactions)
53/// to learn more about RocksDB TransactionDB.
54///
55/// The default thread mode for [`TransactionDB`] is [`SingleThreaded`]
56/// if feature `multi-threaded-cf` is not enabled.
57///
58/// ```
59/// use rocksdb::{DB, Options, TransactionDB, SingleThreaded};
60/// let tempdir = tempfile::Builder::new()
61///     .prefix("_path_for_transaction_db")
62///     .tempdir()
63///     .expect("Failed to create temporary path for the _path_for_transaction_db");
64/// let path = tempdir.path();
65/// {
66///     let db: TransactionDB = TransactionDB::open_default(path).unwrap();
67///     db.put(b"my key", b"my value").unwrap();
68///
69///     // create transaction
70///     let txn = db.transaction();
71///     txn.put(b"key2", b"value2");
72///     txn.put(b"key3", b"value3");
73///     txn.commit().unwrap();
74/// }
75/// let _ = DB::destroy(&Options::default(), path);
76/// ```
77///
78/// [`SingleThreaded`]: crate::SingleThreaded
79pub struct TransactionDB<T: ThreadMode = DefaultThreadMode> {
80    pub(crate) inner: *mut ffi::rocksdb_transactiondb_t,
81    cfs: T,
82    path: PathBuf,
83    // prepared 2pc transactions.
84    prepared: Mutex<Vec<*mut rocksdb_transaction_t>>,
85    _outlive: Vec<OptionsMustOutliveDB>,
86}
87
88unsafe impl<T: ThreadMode> Send for TransactionDB<T> {}
89unsafe impl<T: ThreadMode> Sync for TransactionDB<T> {}
90
91impl<T: ThreadMode> DBAccess for TransactionDB<T> {
92    unsafe fn create_snapshot(&self) -> *const ffi::rocksdb_snapshot_t {
93        ffi::rocksdb_transactiondb_create_snapshot(self.inner)
94    }
95
96    unsafe fn release_snapshot(&self, snapshot: *const ffi::rocksdb_snapshot_t) {
97        ffi::rocksdb_transactiondb_release_snapshot(self.inner, snapshot);
98    }
99
100    unsafe fn create_iterator(&self, readopts: &ReadOptions) -> *mut ffi::rocksdb_iterator_t {
101        ffi::rocksdb_transactiondb_create_iterator(self.inner, readopts.inner)
102    }
103
104    unsafe fn create_iterator_cf(
105        &self,
106        cf_handle: *mut ffi::rocksdb_column_family_handle_t,
107        readopts: &ReadOptions,
108    ) -> *mut ffi::rocksdb_iterator_t {
109        ffi::rocksdb_transactiondb_create_iterator_cf(self.inner, readopts.inner, cf_handle)
110    }
111
112    fn get_opt<K: AsRef<[u8]>>(
113        &self,
114        key: K,
115        readopts: &ReadOptions,
116    ) -> Result<Option<Vec<u8>>, Error> {
117        self.get_opt(key, readopts)
118    }
119
120    fn get_cf_opt<K: AsRef<[u8]>>(
121        &self,
122        cf: &impl AsColumnFamilyRef,
123        key: K,
124        readopts: &ReadOptions,
125    ) -> Result<Option<Vec<u8>>, Error> {
126        self.get_cf_opt(cf, key, readopts)
127    }
128
129    fn get_pinned_opt<K: AsRef<[u8]>>(
130        &self,
131        key: K,
132        readopts: &ReadOptions,
133    ) -> Result<Option<DBPinnableSlice>, Error> {
134        self.get_pinned_opt(key, readopts)
135    }
136
137    fn get_pinned_cf_opt<K: AsRef<[u8]>>(
138        &self,
139        cf: &impl AsColumnFamilyRef,
140        key: K,
141        readopts: &ReadOptions,
142    ) -> Result<Option<DBPinnableSlice>, Error> {
143        self.get_pinned_cf_opt(cf, key, readopts)
144    }
145
146    fn multi_get_opt<K, I>(
147        &self,
148        keys: I,
149        readopts: &ReadOptions,
150    ) -> Vec<Result<Option<Vec<u8>>, Error>>
151    where
152        K: AsRef<[u8]>,
153        I: IntoIterator<Item = K>,
154    {
155        self.multi_get_opt(keys, readopts)
156    }
157
158    fn multi_get_cf_opt<'b, K, I, W>(
159        &self,
160        keys_cf: I,
161        readopts: &ReadOptions,
162    ) -> Vec<Result<Option<Vec<u8>>, Error>>
163    where
164        K: AsRef<[u8]>,
165        I: IntoIterator<Item = (&'b W, K)>,
166        W: AsColumnFamilyRef + 'b,
167    {
168        self.multi_get_cf_opt(keys_cf, readopts)
169    }
170}
171
172impl<T: ThreadMode> TransactionDB<T> {
173    /// Opens a database with default options.
174    pub fn open_default<P: AsRef<Path>>(path: P) -> Result<Self, Error> {
175        let mut opts = Options::default();
176        opts.create_if_missing(true);
177        let txn_db_opts = TransactionDBOptions::default();
178        Self::open(&opts, &txn_db_opts, path)
179    }
180
181    /// Opens the database with the specified options.
182    pub fn open<P: AsRef<Path>>(
183        opts: &Options,
184        txn_db_opts: &TransactionDBOptions,
185        path: P,
186    ) -> Result<Self, Error> {
187        Self::open_cf(opts, txn_db_opts, path, None::<&str>)
188    }
189
190    /// Opens a database with the given database options and column family names.
191    ///
192    /// Column families opened using this function will be created with default `Options`.
193    pub fn open_cf<P, I, N>(
194        opts: &Options,
195        txn_db_opts: &TransactionDBOptions,
196        path: P,
197        cfs: I,
198    ) -> Result<Self, Error>
199    where
200        P: AsRef<Path>,
201        I: IntoIterator<Item = N>,
202        N: AsRef<str>,
203    {
204        let cfs = cfs
205            .into_iter()
206            .map(|name| ColumnFamilyDescriptor::new(name.as_ref(), Options::default()));
207
208        Self::open_cf_descriptors_internal(opts, txn_db_opts, path, cfs)
209    }
210
211    /// Opens a database with the given database options and column family descriptors.
212    pub fn open_cf_descriptors<P, I>(
213        opts: &Options,
214        txn_db_opts: &TransactionDBOptions,
215        path: P,
216        cfs: I,
217    ) -> Result<Self, Error>
218    where
219        P: AsRef<Path>,
220        I: IntoIterator<Item = ColumnFamilyDescriptor>,
221    {
222        Self::open_cf_descriptors_internal(opts, txn_db_opts, path, cfs)
223    }
224
225    /// Internal implementation for opening RocksDB.
226    fn open_cf_descriptors_internal<P, I>(
227        opts: &Options,
228        txn_db_opts: &TransactionDBOptions,
229        path: P,
230        cfs: I,
231    ) -> Result<Self, Error>
232    where
233        P: AsRef<Path>,
234        I: IntoIterator<Item = ColumnFamilyDescriptor>,
235    {
236        let cfs: Vec<_> = cfs.into_iter().collect();
237        let outlive = iter::once(opts.outlive.clone())
238            .chain(cfs.iter().map(|cf| cf.options.outlive.clone()))
239            .collect();
240
241        let cpath = to_cpath(&path)?;
242
243        if let Err(e) = fs::create_dir_all(&path) {
244            return Err(Error::new(format!(
245                "Failed to create RocksDB directory: `{e:?}`."
246            )));
247        }
248
249        let db: *mut ffi::rocksdb_transactiondb_t;
250        let mut cf_map = BTreeMap::new();
251
252        if cfs.is_empty() {
253            db = Self::open_raw(opts, txn_db_opts, &cpath)?;
254        } else {
255            let mut cfs_v = cfs;
256            // Always open the default column family.
257            if !cfs_v.iter().any(|cf| cf.name == DEFAULT_COLUMN_FAMILY_NAME) {
258                cfs_v.push(ColumnFamilyDescriptor {
259                    name: String::from(DEFAULT_COLUMN_FAMILY_NAME),
260                    options: Options::default(),
261                    ttl: ColumnFamilyTtl::SameAsDb, // it will have ttl specified in `DBWithThreadMode::open_with_ttl`
262                });
263            }
264            // We need to store our CStrings in an intermediate vector
265            // so that their pointers remain valid.
266            let c_cfs: Vec<CString> = cfs_v
267                .iter()
268                .map(|cf| CString::new(cf.name.as_bytes()).unwrap())
269                .collect();
270
271            let cfnames: Vec<_> = c_cfs.iter().map(|cf| cf.as_ptr()).collect();
272
273            // These handles will be populated by DB.
274            let mut cfhandles: Vec<_> = cfs_v.iter().map(|_| ptr::null_mut()).collect();
275
276            let cfopts: Vec<_> = cfs_v
277                .iter()
278                .map(|cf| cf.options.inner.cast_const())
279                .collect();
280
281            db = Self::open_cf_raw(
282                opts,
283                txn_db_opts,
284                &cpath,
285                &cfs_v,
286                &cfnames,
287                &cfopts,
288                &mut cfhandles,
289            )?;
290
291            for handle in &cfhandles {
292                if handle.is_null() {
293                    return Err(Error::new(
294                        "Received null column family handle from DB.".to_owned(),
295                    ));
296                }
297            }
298
299            for (cf_desc, inner) in cfs_v.iter().zip(cfhandles) {
300                cf_map.insert(cf_desc.name.clone(), inner);
301            }
302        }
303
304        if db.is_null() {
305            return Err(Error::new("Could not initialize database.".to_owned()));
306        }
307
308        let prepared = unsafe {
309            let mut cnt = 0;
310            let ptr = ffi::rocksdb_transactiondb_get_prepared_transactions(db, &raw mut cnt);
311            let mut vec = vec![std::ptr::null_mut(); cnt];
312            if !ptr.is_null() {
313                std::ptr::copy_nonoverlapping(ptr, vec.as_mut_ptr(), cnt);
314                ffi::rocksdb_free(ptr as *mut c_void);
315            }
316            vec
317        };
318
319        Ok(TransactionDB {
320            inner: db,
321            cfs: T::new_cf_map_internal(cf_map),
322            path: path.as_ref().to_path_buf(),
323            prepared: Mutex::new(prepared),
324            _outlive: outlive,
325        })
326    }
327
328    fn open_raw(
329        opts: &Options,
330        txn_db_opts: &TransactionDBOptions,
331        cpath: &CString,
332    ) -> Result<*mut ffi::rocksdb_transactiondb_t, Error> {
333        unsafe {
334            let db = ffi_try!(ffi::rocksdb_transactiondb_open(
335                opts.inner,
336                txn_db_opts.inner,
337                cpath.as_ptr()
338            ));
339            Ok(db)
340        }
341    }
342
343    fn open_cf_raw(
344        opts: &Options,
345        txn_db_opts: &TransactionDBOptions,
346        cpath: &CString,
347        cfs_v: &[ColumnFamilyDescriptor],
348        cfnames: &[*const c_char],
349        cfopts: &[*const ffi::rocksdb_options_t],
350        cfhandles: &mut [*mut ffi::rocksdb_column_family_handle_t],
351    ) -> Result<*mut ffi::rocksdb_transactiondb_t, Error> {
352        unsafe {
353            let db = ffi_try!(ffi::rocksdb_transactiondb_open_column_families(
354                opts.inner,
355                txn_db_opts.inner,
356                cpath.as_ptr(),
357                cfs_v.len() as c_int,
358                cfnames.as_ptr(),
359                cfopts.as_ptr(),
360                cfhandles.as_mut_ptr(),
361            ));
362            Ok(db)
363        }
364    }
365
366    fn create_inner_cf_handle(
367        &self,
368        name: &str,
369        opts: &Options,
370    ) -> Result<*mut ffi::rocksdb_column_family_handle_t, Error> {
371        let cf_name = CString::new(name.as_bytes()).map_err(|_| {
372            Error::new("Failed to convert path to CString when creating cf".to_owned())
373        })?;
374
375        Ok(unsafe {
376            ffi_try!(ffi::rocksdb_transactiondb_create_column_family(
377                self.inner,
378                opts.inner,
379                cf_name.as_ptr(),
380            ))
381        })
382    }
383
384    pub fn list_cf<P: AsRef<Path>>(opts: &Options, path: P) -> Result<Vec<String>, Error> {
385        DB::list_cf(opts, path)
386    }
387
388    pub fn destroy<P: AsRef<Path>>(opts: &Options, path: P) -> Result<(), Error> {
389        DB::destroy(opts, path)
390    }
391
392    pub fn repair<P: AsRef<Path>>(opts: &Options, path: P) -> Result<(), Error> {
393        DB::repair(opts, path)
394    }
395
396    pub fn path(&self) -> &Path {
397        self.path.as_path()
398    }
399
400    /// Flushes the WAL buffer. If `sync` is set to `true`, also syncs
401    /// the data to disk.
402    pub fn flush_wal(&self, sync: bool) -> Result<(), Error> {
403        unsafe {
404            ffi_try!(ffi::rocksdb_transactiondb_flush_wal(
405                self.inner,
406                c_uchar::from(sync)
407            ));
408        }
409        Ok(())
410    }
411
412    /// Flushes database memtables to SST files on the disk.
413    pub fn flush_opt(&self, flushopts: &FlushOptions) -> Result<(), Error> {
414        unsafe {
415            ffi_try!(ffi::rocksdb_transactiondb_flush(
416                self.inner,
417                flushopts.inner
418            ));
419        }
420        Ok(())
421    }
422
423    /// Flushes database memtables to SST files on the disk using default options.
424    pub fn flush(&self) -> Result<(), Error> {
425        self.flush_opt(&FlushOptions::default())
426    }
427
428    /// Flushes database memtables to SST files on the disk for a given column family.
429    pub fn flush_cf_opt(
430        &self,
431        cf: &impl AsColumnFamilyRef,
432        flushopts: &FlushOptions,
433    ) -> Result<(), Error> {
434        unsafe {
435            ffi_try!(ffi::rocksdb_transactiondb_flush_cf(
436                self.inner,
437                flushopts.inner,
438                cf.inner()
439            ));
440        }
441        Ok(())
442    }
443
444    /// Flushes multiple column families.
445    ///
446    /// If atomic flush is not enabled, it is equivalent to calling flush_cf multiple times.
447    /// If atomic flush is enabled, it will flush all column families specified in `cfs` up to the latest sequence
448    /// number at the time when flush is requested.
449    pub fn flush_cfs_opt(
450        &self,
451        cfs: &[&impl AsColumnFamilyRef],
452        opts: &FlushOptions,
453    ) -> Result<(), Error> {
454        let mut cfs = cfs.iter().map(|cf| cf.inner()).collect::<Vec<_>>();
455        unsafe {
456            ffi_try!(ffi::rocksdb_transactiondb_flush_cfs(
457                self.inner,
458                opts.inner,
459                cfs.as_mut_ptr(),
460                cfs.len() as c_int,
461            ));
462        }
463        Ok(())
464    }
465
466    /// Flushes database memtables to SST files on the disk for a given column family using default
467    /// options.
468    pub fn flush_cf(&self, cf: &impl AsColumnFamilyRef) -> Result<(), Error> {
469        self.flush_cf_opt(cf, &FlushOptions::default())
470    }
471
472    /// Creates a transaction with default options.
473    pub fn transaction(&self) -> Transaction<Self> {
474        self.transaction_opt(&WriteOptions::default(), &TransactionOptions::default())
475    }
476
477    /// Creates a transaction with options.
478    pub fn transaction_opt<'a>(
479        &'a self,
480        write_opts: &WriteOptions,
481        txn_opts: &TransactionOptions,
482    ) -> Transaction<'a, Self> {
483        Transaction {
484            inner: unsafe {
485                ffi::rocksdb_transaction_begin(
486                    self.inner,
487                    write_opts.inner,
488                    txn_opts.inner,
489                    std::ptr::null_mut(),
490                )
491            },
492            _marker: PhantomData,
493        }
494    }
495
496    /// Get all prepared transactions for recovery.
497    ///
498    /// This function is expected to call once after open database.
499    /// User should commit or rollback all transactions before start other transactions.
500    pub fn prepared_transactions(&self) -> Vec<Transaction<Self>> {
501        self.prepared
502            .lock()
503            .unwrap()
504            .drain(0..)
505            .map(|inner| Transaction {
506                inner,
507                _marker: PhantomData,
508            })
509            .collect()
510    }
511
512    /// Returns the bytes associated with a key value.
513    pub fn get<K: AsRef<[u8]>>(&self, key: K) -> Result<Option<Vec<u8>>, Error> {
514        self.get_pinned(key).map(|x| x.map(|v| v.as_ref().to_vec()))
515    }
516
517    /// Returns the bytes associated with a key value and the given column family.
518    pub fn get_cf<K: AsRef<[u8]>>(
519        &self,
520        cf: &impl AsColumnFamilyRef,
521        key: K,
522    ) -> Result<Option<Vec<u8>>, Error> {
523        self.get_pinned_cf(cf, key)
524            .map(|x| x.map(|v| v.as_ref().to_vec()))
525    }
526
527    /// Returns the bytes associated with a key value with read options.
528    pub fn get_opt<K: AsRef<[u8]>>(
529        &self,
530        key: K,
531        readopts: &ReadOptions,
532    ) -> Result<Option<Vec<u8>>, Error> {
533        self.get_pinned_opt(key, readopts)
534            .map(|x| x.map(|v| v.as_ref().to_vec()))
535    }
536
537    /// Returns the bytes associated with a key value and the given column family with read options.
538    pub fn get_cf_opt<K: AsRef<[u8]>>(
539        &self,
540        cf: &impl AsColumnFamilyRef,
541        key: K,
542        readopts: &ReadOptions,
543    ) -> Result<Option<Vec<u8>>, Error> {
544        self.get_pinned_cf_opt(cf, key, readopts)
545            .map(|x| x.map(|v| v.as_ref().to_vec()))
546    }
547
548    pub fn get_pinned<K: AsRef<[u8]>>(&self, key: K) -> Result<Option<DBPinnableSlice>, Error> {
549        self.get_pinned_opt(key, &ReadOptions::default())
550    }
551
552    /// Returns the bytes associated with a key value and the given column family.
553    pub fn get_pinned_cf<K: AsRef<[u8]>>(
554        &self,
555        cf: &impl AsColumnFamilyRef,
556        key: K,
557    ) -> Result<Option<DBPinnableSlice>, Error> {
558        self.get_pinned_cf_opt(cf, key, &ReadOptions::default())
559    }
560
561    /// Returns the bytes associated with a key value with read options.
562    pub fn get_pinned_opt<K: AsRef<[u8]>>(
563        &self,
564        key: K,
565        readopts: &ReadOptions,
566    ) -> Result<Option<DBPinnableSlice>, Error> {
567        let key = key.as_ref();
568        unsafe {
569            let val = ffi_try!(ffi::rocksdb_transactiondb_get_pinned(
570                self.inner,
571                readopts.inner,
572                key.as_ptr() as *const c_char,
573                key.len() as size_t,
574            ));
575            if val.is_null() {
576                Ok(None)
577            } else {
578                Ok(Some(DBPinnableSlice::from_c(val)))
579            }
580        }
581    }
582
583    /// Returns the bytes associated with a key value and the given column family with read options.
584    pub fn get_pinned_cf_opt<K: AsRef<[u8]>>(
585        &self,
586        cf: &impl AsColumnFamilyRef,
587        key: K,
588        readopts: &ReadOptions,
589    ) -> Result<Option<DBPinnableSlice>, Error> {
590        let key = key.as_ref();
591        unsafe {
592            let val = ffi_try!(ffi::rocksdb_transactiondb_get_pinned_cf(
593                self.inner,
594                readopts.inner,
595                cf.inner(),
596                key.as_ptr() as *const c_char,
597                key.len() as size_t,
598            ));
599            if val.is_null() {
600                Ok(None)
601            } else {
602                Ok(Some(DBPinnableSlice::from_c(val)))
603            }
604        }
605    }
606
607    /// Return the values associated with the given keys.
608    pub fn multi_get<K, I>(&self, keys: I) -> Vec<Result<Option<Vec<u8>>, Error>>
609    where
610        K: AsRef<[u8]>,
611        I: IntoIterator<Item = K>,
612    {
613        self.multi_get_opt(keys, &ReadOptions::default())
614    }
615
616    /// Return the values associated with the given keys using read options.
617    pub fn multi_get_opt<K, I>(
618        &self,
619        keys: I,
620        readopts: &ReadOptions,
621    ) -> Vec<Result<Option<Vec<u8>>, Error>>
622    where
623        K: AsRef<[u8]>,
624        I: IntoIterator<Item = K>,
625    {
626        let (keys, keys_sizes): (Vec<Box<[u8]>>, Vec<_>) = keys
627            .into_iter()
628            .map(|key| {
629                let key = key.as_ref();
630                (Box::from(key), key.len())
631            })
632            .unzip();
633        let ptr_keys: Vec<_> = keys.iter().map(|k| k.as_ptr() as *const c_char).collect();
634
635        let mut values = vec![ptr::null_mut(); keys.len()];
636        let mut values_sizes = vec![0_usize; keys.len()];
637        let mut errors = vec![ptr::null_mut(); keys.len()];
638        unsafe {
639            ffi::rocksdb_transactiondb_multi_get(
640                self.inner,
641                readopts.inner,
642                ptr_keys.len(),
643                ptr_keys.as_ptr(),
644                keys_sizes.as_ptr(),
645                values.as_mut_ptr(),
646                values_sizes.as_mut_ptr(),
647                errors.as_mut_ptr(),
648            );
649        }
650
651        convert_values(values, values_sizes, errors)
652    }
653
654    /// Return the values associated with the given keys and column families.
655    pub fn multi_get_cf<'a, 'b: 'a, K, I, W>(
656        &'a self,
657        keys: I,
658    ) -> Vec<Result<Option<Vec<u8>>, Error>>
659    where
660        K: AsRef<[u8]>,
661        I: IntoIterator<Item = (&'b W, K)>,
662        W: 'b + AsColumnFamilyRef,
663    {
664        self.multi_get_cf_opt(keys, &ReadOptions::default())
665    }
666
667    /// Return the values associated with the given keys and column families using read options.
668    pub fn multi_get_cf_opt<'a, 'b: 'a, K, I, W>(
669        &'a self,
670        keys: I,
671        readopts: &ReadOptions,
672    ) -> Vec<Result<Option<Vec<u8>>, Error>>
673    where
674        K: AsRef<[u8]>,
675        I: IntoIterator<Item = (&'b W, K)>,
676        W: 'b + AsColumnFamilyRef,
677    {
678        let (cfs_and_keys, keys_sizes): (Vec<(_, Box<[u8]>)>, Vec<_>) = keys
679            .into_iter()
680            .map(|(cf, key)| {
681                let key = key.as_ref();
682                ((cf, Box::from(key)), key.len())
683            })
684            .unzip();
685        let ptr_keys: Vec<_> = cfs_and_keys
686            .iter()
687            .map(|(_, k)| k.as_ptr() as *const c_char)
688            .collect();
689        let ptr_cfs: Vec<_> = cfs_and_keys
690            .iter()
691            .map(|(c, _)| c.inner().cast_const())
692            .collect();
693
694        let mut values = vec![ptr::null_mut(); ptr_keys.len()];
695        let mut values_sizes = vec![0_usize; ptr_keys.len()];
696        let mut errors = vec![ptr::null_mut(); ptr_keys.len()];
697        unsafe {
698            ffi::rocksdb_transactiondb_multi_get_cf(
699                self.inner,
700                readopts.inner,
701                ptr_cfs.as_ptr(),
702                ptr_keys.len(),
703                ptr_keys.as_ptr(),
704                keys_sizes.as_ptr(),
705                values.as_mut_ptr(),
706                values_sizes.as_mut_ptr(),
707                errors.as_mut_ptr(),
708            );
709        }
710
711        convert_values(values, values_sizes, errors)
712    }
713
714    pub fn put<K, V>(&self, key: K, value: V) -> Result<(), Error>
715    where
716        K: AsRef<[u8]>,
717        V: AsRef<[u8]>,
718    {
719        self.put_opt(key, value, &WriteOptions::default())
720    }
721
722    pub fn put_cf<K, V>(&self, cf: &impl AsColumnFamilyRef, key: K, value: V) -> Result<(), Error>
723    where
724        K: AsRef<[u8]>,
725        V: AsRef<[u8]>,
726    {
727        self.put_cf_opt(cf, key, value, &WriteOptions::default())
728    }
729
730    pub fn put_opt<K, V>(&self, key: K, value: V, writeopts: &WriteOptions) -> Result<(), Error>
731    where
732        K: AsRef<[u8]>,
733        V: AsRef<[u8]>,
734    {
735        let key = key.as_ref();
736        let value = value.as_ref();
737        unsafe {
738            ffi_try!(ffi::rocksdb_transactiondb_put(
739                self.inner,
740                writeopts.inner,
741                key.as_ptr() as *const c_char,
742                key.len() as size_t,
743                value.as_ptr() as *const c_char,
744                value.len() as size_t
745            ));
746        }
747        Ok(())
748    }
749
750    pub fn put_cf_opt<K, V>(
751        &self,
752        cf: &impl AsColumnFamilyRef,
753        key: K,
754        value: V,
755        writeopts: &WriteOptions,
756    ) -> Result<(), Error>
757    where
758        K: AsRef<[u8]>,
759        V: AsRef<[u8]>,
760    {
761        let key = key.as_ref();
762        let value = value.as_ref();
763        unsafe {
764            ffi_try!(ffi::rocksdb_transactiondb_put_cf(
765                self.inner,
766                writeopts.inner,
767                cf.inner(),
768                key.as_ptr() as *const c_char,
769                key.len() as size_t,
770                value.as_ptr() as *const c_char,
771                value.len() as size_t
772            ));
773        }
774        Ok(())
775    }
776
777    pub fn write(&self, batch: WriteBatchWithTransaction<true>) -> Result<(), Error> {
778        self.write_opt(batch, &WriteOptions::default())
779    }
780
781    pub fn write_opt(
782        &self,
783        batch: WriteBatchWithTransaction<true>,
784        writeopts: &WriteOptions,
785    ) -> Result<(), Error> {
786        unsafe {
787            ffi_try!(ffi::rocksdb_transactiondb_write(
788                self.inner,
789                writeopts.inner,
790                batch.inner
791            ));
792        }
793        Ok(())
794    }
795
796    pub fn merge<K, V>(&self, key: K, value: V) -> Result<(), Error>
797    where
798        K: AsRef<[u8]>,
799        V: AsRef<[u8]>,
800    {
801        self.merge_opt(key, value, &WriteOptions::default())
802    }
803
804    pub fn merge_cf<K, V>(&self, cf: &impl AsColumnFamilyRef, key: K, value: V) -> Result<(), Error>
805    where
806        K: AsRef<[u8]>,
807        V: AsRef<[u8]>,
808    {
809        self.merge_cf_opt(cf, key, value, &WriteOptions::default())
810    }
811
812    pub fn merge_opt<K, V>(&self, key: K, value: V, writeopts: &WriteOptions) -> Result<(), Error>
813    where
814        K: AsRef<[u8]>,
815        V: AsRef<[u8]>,
816    {
817        let key = key.as_ref();
818        let value = value.as_ref();
819        unsafe {
820            ffi_try!(ffi::rocksdb_transactiondb_merge(
821                self.inner,
822                writeopts.inner,
823                key.as_ptr() as *const c_char,
824                key.len() as size_t,
825                value.as_ptr() as *const c_char,
826                value.len() as size_t,
827            ));
828            Ok(())
829        }
830    }
831
832    pub fn merge_cf_opt<K, V>(
833        &self,
834        cf: &impl AsColumnFamilyRef,
835        key: K,
836        value: V,
837        writeopts: &WriteOptions,
838    ) -> Result<(), Error>
839    where
840        K: AsRef<[u8]>,
841        V: AsRef<[u8]>,
842    {
843        let key = key.as_ref();
844        let value = value.as_ref();
845        unsafe {
846            ffi_try!(ffi::rocksdb_transactiondb_merge_cf(
847                self.inner,
848                writeopts.inner,
849                cf.inner(),
850                key.as_ptr() as *const c_char,
851                key.len() as size_t,
852                value.as_ptr() as *const c_char,
853                value.len() as size_t,
854            ));
855            Ok(())
856        }
857    }
858
859    pub fn delete<K: AsRef<[u8]>>(&self, key: K) -> Result<(), Error> {
860        self.delete_opt(key, &WriteOptions::default())
861    }
862
863    pub fn delete_cf<K: AsRef<[u8]>>(
864        &self,
865        cf: &impl AsColumnFamilyRef,
866        key: K,
867    ) -> Result<(), Error> {
868        self.delete_cf_opt(cf, key, &WriteOptions::default())
869    }
870
871    pub fn delete_opt<K: AsRef<[u8]>>(
872        &self,
873        key: K,
874        writeopts: &WriteOptions,
875    ) -> Result<(), Error> {
876        let key = key.as_ref();
877        unsafe {
878            ffi_try!(ffi::rocksdb_transactiondb_delete(
879                self.inner,
880                writeopts.inner,
881                key.as_ptr() as *const c_char,
882                key.len() as size_t,
883            ));
884        }
885        Ok(())
886    }
887
888    pub fn delete_cf_opt<K: AsRef<[u8]>>(
889        &self,
890        cf: &impl AsColumnFamilyRef,
891        key: K,
892        writeopts: &WriteOptions,
893    ) -> Result<(), Error> {
894        let key = key.as_ref();
895        unsafe {
896            ffi_try!(ffi::rocksdb_transactiondb_delete_cf(
897                self.inner,
898                writeopts.inner,
899                cf.inner(),
900                key.as_ptr() as *const c_char,
901                key.len() as size_t,
902            ));
903        }
904        Ok(())
905    }
906
907    pub fn iterator<'a: 'b, 'b>(
908        &'a self,
909        mode: IteratorMode,
910    ) -> DBIteratorWithThreadMode<'b, Self> {
911        let readopts = ReadOptions::default();
912        self.iterator_opt(mode, readopts)
913    }
914
915    pub fn iterator_opt<'a: 'b, 'b>(
916        &'a self,
917        mode: IteratorMode,
918        readopts: ReadOptions,
919    ) -> DBIteratorWithThreadMode<'b, Self> {
920        DBIteratorWithThreadMode::new(self, readopts, mode)
921    }
922
923    /// Opens an iterator using the provided ReadOptions.
924    /// This is used when you want to iterate over a specific ColumnFamily with a modified ReadOptions
925    pub fn iterator_cf_opt<'a: 'b, 'b>(
926        &'a self,
927        cf_handle: &impl AsColumnFamilyRef,
928        readopts: ReadOptions,
929        mode: IteratorMode,
930    ) -> DBIteratorWithThreadMode<'b, Self> {
931        DBIteratorWithThreadMode::new_cf(self, cf_handle.inner(), readopts, mode)
932    }
933
934    /// Opens an iterator with `set_total_order_seek` enabled.
935    /// This must be used to iterate across prefixes when `set_memtable_factory` has been called
936    /// with a Hash-based implementation.
937    pub fn full_iterator<'a: 'b, 'b>(
938        &'a self,
939        mode: IteratorMode,
940    ) -> DBIteratorWithThreadMode<'b, Self> {
941        let mut opts = ReadOptions::default();
942        opts.set_total_order_seek(true);
943        DBIteratorWithThreadMode::new(self, opts, mode)
944    }
945
946    pub fn prefix_iterator<'a: 'b, 'b, P: AsRef<[u8]>>(
947        &'a self,
948        prefix: P,
949    ) -> DBIteratorWithThreadMode<'b, Self> {
950        let mut opts = ReadOptions::default();
951        opts.set_prefix_same_as_start(true);
952        DBIteratorWithThreadMode::new(
953            self,
954            opts,
955            IteratorMode::From(prefix.as_ref(), Direction::Forward),
956        )
957    }
958
959    pub fn iterator_cf<'a: 'b, 'b>(
960        &'a self,
961        cf_handle: &impl AsColumnFamilyRef,
962        mode: IteratorMode,
963    ) -> DBIteratorWithThreadMode<'b, Self> {
964        let opts = ReadOptions::default();
965        DBIteratorWithThreadMode::new_cf(self, cf_handle.inner(), opts, mode)
966    }
967
968    pub fn full_iterator_cf<'a: 'b, 'b>(
969        &'a self,
970        cf_handle: &impl AsColumnFamilyRef,
971        mode: IteratorMode,
972    ) -> DBIteratorWithThreadMode<'b, Self> {
973        let mut opts = ReadOptions::default();
974        opts.set_total_order_seek(true);
975        DBIteratorWithThreadMode::new_cf(self, cf_handle.inner(), opts, mode)
976    }
977
978    pub fn prefix_iterator_cf<'a, P: AsRef<[u8]>>(
979        &'a self,
980        cf_handle: &impl AsColumnFamilyRef,
981        prefix: P,
982    ) -> DBIteratorWithThreadMode<'a, Self> {
983        let mut opts = ReadOptions::default();
984        opts.set_prefix_same_as_start(true);
985        DBIteratorWithThreadMode::<'a, Self>::new_cf(
986            self,
987            cf_handle.inner(),
988            opts,
989            IteratorMode::From(prefix.as_ref(), Direction::Forward),
990        )
991    }
992
993    /// Opens a raw iterator over the database, using the default read options
994    pub fn raw_iterator<'a: 'b, 'b>(&'a self) -> DBRawIteratorWithThreadMode<'b, Self> {
995        let opts = ReadOptions::default();
996        DBRawIteratorWithThreadMode::new(self, opts)
997    }
998
999    /// Opens a raw iterator over the given column family, using the default read options
1000    pub fn raw_iterator_cf<'a: 'b, 'b>(
1001        &'a self,
1002        cf_handle: &impl AsColumnFamilyRef,
1003    ) -> DBRawIteratorWithThreadMode<'b, Self> {
1004        let opts = ReadOptions::default();
1005        DBRawIteratorWithThreadMode::new_cf(self, cf_handle.inner(), opts)
1006    }
1007
1008    /// Opens a raw iterator over the database, using the given read options
1009    pub fn raw_iterator_opt<'a: 'b, 'b>(
1010        &'a self,
1011        readopts: ReadOptions,
1012    ) -> DBRawIteratorWithThreadMode<'b, Self> {
1013        DBRawIteratorWithThreadMode::new(self, readopts)
1014    }
1015
1016    /// Opens a raw iterator over the given column family, using the given read options
1017    pub fn raw_iterator_cf_opt<'a: 'b, 'b>(
1018        &'a self,
1019        cf_handle: &impl AsColumnFamilyRef,
1020        readopts: ReadOptions,
1021    ) -> DBRawIteratorWithThreadMode<'b, Self> {
1022        DBRawIteratorWithThreadMode::new_cf(self, cf_handle.inner(), readopts)
1023    }
1024
1025    pub fn snapshot(&self) -> SnapshotWithThreadMode<Self> {
1026        SnapshotWithThreadMode::<Self>::new(self)
1027    }
1028
1029    fn drop_column_family<C>(
1030        &self,
1031        cf_inner: *mut ffi::rocksdb_column_family_handle_t,
1032        _cf: C,
1033    ) -> Result<(), Error> {
1034        unsafe {
1035            // first mark the column family as dropped
1036            ffi_try!(ffi::rocksdb_drop_column_family(
1037                self.inner as *mut ffi::rocksdb_t,
1038                cf_inner
1039            ));
1040        }
1041        // Since `_cf` is dropped here, the column family handle is destroyed
1042        // and any resources (mem, files) are reclaimed.
1043        Ok(())
1044    }
1045}
1046
1047impl TransactionDB<SingleThreaded> {
1048    /// Creates column family with given name and options.
1049    pub fn create_cf<N: AsRef<str>>(&mut self, name: N, opts: &Options) -> Result<(), Error> {
1050        let inner = self.create_inner_cf_handle(name.as_ref(), opts)?;
1051        self.cfs
1052            .cfs
1053            .insert(name.as_ref().to_string(), ColumnFamily { inner });
1054        Ok(())
1055    }
1056
1057    /// Returns the underlying column family handle.
1058    pub fn cf_handle(&self, name: &str) -> Option<&ColumnFamily> {
1059        self.cfs.cfs.get(name)
1060    }
1061
1062    /// Drops the column family with the given name
1063    pub fn drop_cf(&mut self, name: &str) -> Result<(), Error> {
1064        if let Some(cf) = self.cfs.cfs.remove(name) {
1065            self.drop_column_family(cf.inner, cf)
1066        } else {
1067            Err(Error::new(format!("Invalid column family: {name}")))
1068        }
1069    }
1070}
1071
1072impl TransactionDB<MultiThreaded> {
1073    /// Creates column family with given name and options.
1074    pub fn create_cf<N: AsRef<str>>(&self, name: N, opts: &Options) -> Result<(), Error> {
1075        // Note that we acquire the cfs lock before inserting: otherwise we might race
1076        // another caller who observed the handle as missing.
1077        let mut cfs = self.cfs.cfs.write().unwrap();
1078        let inner = self.create_inner_cf_handle(name.as_ref(), opts)?;
1079        cfs.insert(
1080            name.as_ref().to_string(),
1081            Arc::new(UnboundColumnFamily { inner }),
1082        );
1083        Ok(())
1084    }
1085
1086    /// Returns the underlying column family handle.
1087    pub fn cf_handle(&self, name: &str) -> Option<Arc<BoundColumnFamily>> {
1088        self.cfs
1089            .cfs
1090            .read()
1091            .unwrap()
1092            .get(name)
1093            .cloned()
1094            .map(UnboundColumnFamily::bound_column_family)
1095    }
1096
1097    /// Drops the column family with the given name by internally locking the inner column
1098    /// family map. This avoids needing `&mut self` reference
1099    pub fn drop_cf(&self, name: &str) -> Result<(), Error> {
1100        if let Some(cf) = self.cfs.cfs.write().unwrap().remove(name) {
1101            self.drop_column_family(cf.inner, cf)
1102        } else {
1103            Err(Error::new(format!("Invalid column family: {name}")))
1104        }
1105    }
1106
1107    /// Implementation for property_value et al methods.
1108    ///
1109    /// `name` is the name of the property.  It will be converted into a CString
1110    /// and passed to `get_property` as argument.  `get_property` reads the
1111    /// specified property and either returns NULL or a pointer to a C allocated
1112    /// string; this method takes ownership of that string and will free it at
1113    /// the end. That string is parsed using `parse` callback which produces
1114    /// the returned result.
1115    fn property_value_impl<R>(
1116        name: impl CStrLike,
1117        get_property: impl FnOnce(*const c_char) -> *mut c_char,
1118        parse: impl FnOnce(&str) -> Result<R, Error>,
1119    ) -> Result<Option<R>, Error> {
1120        let value = match name.bake() {
1121            Ok(prop_name) => get_property(prop_name.as_ptr()),
1122            Err(e) => {
1123                return Err(Error::new(format!(
1124                    "Failed to convert property name to CString: {e}"
1125                )));
1126            }
1127        };
1128        if value.is_null() {
1129            return Ok(None);
1130        }
1131        let result = match unsafe { CStr::from_ptr(value) }.to_str() {
1132            Ok(s) => parse(s).map(|value| Some(value)),
1133            Err(e) => Err(Error::new(format!(
1134                "Failed to convert property value to string: {e}"
1135            ))),
1136        };
1137        unsafe {
1138            ffi::rocksdb_free(value as *mut c_void);
1139        }
1140        result
1141    }
1142
1143    /// Retrieves a RocksDB property by name.
1144    ///
1145    /// Full list of properties could be find
1146    /// [here](https://github.com/facebook/rocksdb/blob/08809f5e6cd9cc4bc3958dd4d59457ae78c76660/include/rocksdb/db.h#L428-L634).
1147    pub fn property_value(&self, name: impl CStrLike) -> Result<Option<String>, Error> {
1148        Self::property_value_impl(
1149            name,
1150            |prop_name| unsafe { ffi::rocksdb_transactiondb_property_value(self.inner, prop_name) },
1151            |str_value| Ok(str_value.to_owned()),
1152        )
1153    }
1154
1155    fn parse_property_int_value(value: &str) -> Result<u64, Error> {
1156        value.parse::<u64>().map_err(|err| {
1157            Error::new(format!(
1158                "Failed to convert property value {value} to int: {err}"
1159            ))
1160        })
1161    }
1162
1163    /// Retrieves a RocksDB property and casts it to an integer.
1164    ///
1165    /// Full list of properties that return int values could be find
1166    /// [here](https://github.com/facebook/rocksdb/blob/08809f5e6cd9cc4bc3958dd4d59457ae78c76660/include/rocksdb/db.h#L654-L689).
1167    pub fn property_int_value(&self, name: impl CStrLike) -> Result<Option<u64>, Error> {
1168        Self::property_value_impl(
1169            name,
1170            |prop_name| unsafe { ffi::rocksdb_transactiondb_property_value(self.inner, prop_name) },
1171            Self::parse_property_int_value,
1172        )
1173    }
1174}
1175
1176impl<T: ThreadMode> Drop for TransactionDB<T> {
1177    fn drop(&mut self) {
1178        unsafe {
1179            self.prepared_transactions().clear();
1180            self.cfs.drop_all_cfs_internal();
1181            ffi::rocksdb_transactiondb_close(self.inner);
1182        }
1183    }
1184}