1use 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
50pub struct TransactionDB<T: ThreadMode = DefaultThreadMode> {
80 pub(crate) inner: *mut ffi::rocksdb_transactiondb_t,
81 cfs: T,
82 path: PathBuf,
83 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 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 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 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 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 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 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, });
263 }
264 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 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 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 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 pub fn flush(&self) -> Result<(), Error> {
425 self.flush_opt(&FlushOptions::default())
426 }
427
428 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 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 pub fn flush_cf(&self, cf: &impl AsColumnFamilyRef) -> Result<(), Error> {
469 self.flush_cf_opt(cf, &FlushOptions::default())
470 }
471
472 pub fn transaction(&self) -> Transaction<Self> {
474 self.transaction_opt(&WriteOptions::default(), &TransactionOptions::default())
475 }
476
477 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 ffi_try!(ffi::rocksdb_drop_column_family(
1037 self.inner as *mut ffi::rocksdb_t,
1038 cf_inner
1039 ));
1040 }
1041 Ok(())
1044 }
1045}
1046
1047impl TransactionDB<SingleThreaded> {
1048 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 pub fn cf_handle(&self, name: &str) -> Option<&ColumnFamily> {
1059 self.cfs.cfs.get(name)
1060 }
1061
1062 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 pub fn create_cf<N: AsRef<str>>(&self, name: N, opts: &Options) -> Result<(), Error> {
1075 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 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 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 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 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 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}