1use crate::{
22 btree::{commit_overlay::BTreeChangeSet, BTreeIterator, BTreeTable},
23 column::{
24 hash_key, unpack_node_children, unpack_node_data, ColId, Column, HashColumn, IterState,
25 ReindexBatch, ValueIterState,
26 },
27 error::{try_io, Error, Result},
28 hash::IdentityBuildHasher,
29 index::{Address, PlanOutcome},
30 log::{Log, LogAction},
31 multitree::{Children, NewNode, NodeAddress},
32 options::{Options, CURRENT_VERSION},
33 parking_lot::{
34 Condvar, Mutex, MutexGuard, RwLock, RwLockUpgradableReadGuard, RwLockWriteGuard,
35 },
36 stats::StatSummary,
37 ColumnOptions, Key,
38};
39#[cfg(feature = "bytes")]
40use bytes::Bytes;
41use fs2::FileExt;
42use std::{
43 borrow::Borrow,
44 collections::{BTreeMap, HashMap, HashSet, VecDeque},
45 ops::Bound,
46 sync::{
47 atomic::{AtomicBool, AtomicU64, Ordering},
48 Arc, Weak,
49 },
50 thread,
51};
52
53const MAX_COMMIT_QUEUE_BYTES: usize = 16 * 1024 * 1024;
57const MAX_LOG_QUEUE_BYTES: i64 = 128 * 1024 * 1024;
60const MIN_LOG_SIZE_BYTES: u64 = 64 * 1024 * 1024;
62const KEEP_LOGS: usize = 16;
65const MAX_LOG_FILES: usize = 4;
69
70pub type Value = Vec<u8>;
72
73#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
74pub enum RcValue {
75 #[cfg(feature = "arc")]
76 Arc(Arc<Value>),
77 #[cfg(feature = "bytes")]
78 Bytes(Bytes),
79}
80
81impl AsRef<[u8]> for RcValue {
82 fn as_ref(&self) -> &[u8] {
83 match self {
84 #[cfg(feature = "arc")]
85 Self::Arc(arc) => arc.as_ref(),
86 #[cfg(feature = "bytes")]
87 Self::Bytes(bytes) => bytes.as_ref(),
88 }
89 }
90}
91
92impl Borrow<[u8]> for RcValue {
93 fn borrow(&self) -> &[u8] {
94 self.as_ref()
95 }
96}
97
98#[cfg(feature = "arc")]
99impl From<Value> for RcValue {
100 fn from(value: Value) -> Self {
101 Self::Arc(value.into())
102 }
103}
104
105#[cfg(not(feature = "arc"))]
106impl From<Value> for RcValue {
107 fn from(value: Value) -> Self {
108 Self::Bytes(value.into())
109 }
110}
111
112#[cfg(feature = "arc")]
113impl From<Arc<Value>> for RcValue {
114 fn from(value: Arc<Value>) -> Self {
115 Self::Arc(value)
116 }
117}
118
119#[cfg(feature = "bytes")]
120impl From<Bytes> for RcValue {
121 fn from(value: Bytes) -> Self {
122 Self::Bytes(value)
123 }
124}
125
126#[cfg(test)]
127impl<const N: usize> TryFrom<RcValue> for [u8; N] {
128 type Error = <[u8; N] as TryFrom<Vec<u8>>>::Error;
129
130 fn try_from(value: RcValue) -> std::result::Result<Self, Self::Error> {
131 value.as_ref().to_vec().try_into()
132 }
133}
134
135#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
136pub struct RcKey(Arc<Vec<u8>>);
137
138impl AsRef<[u8]> for RcKey {
139 fn as_ref(&self) -> &[u8] {
140 self.0.as_ref()
141 }
142}
143
144impl Borrow<[u8]> for RcKey {
145 fn borrow(&self) -> &[u8] {
146 self.as_ref()
147 }
148}
149
150impl From<Vec<u8>> for RcKey {
151 fn from(key: Vec<u8>) -> Self {
152 Self(key.into())
153 }
154}
155
156#[derive(Debug, Default)]
158struct Commit {
159 id: u64,
162 bytes: usize,
165 changeset: CommitChangeSet,
167}
168
169#[derive(Debug, Default)]
171struct CommitQueue {
172 record_id: u64,
174 bytes: usize,
176 commits: VecDeque<Commit>,
178}
179
180#[derive(Debug)]
181struct Trees {
182 readers: HashMap<Key, Weak<RwLock<Box<dyn TreeReader + Send + Sync>>>, IdentityBuildHasher>,
183 to_dereference: HashMap<Key, usize>,
185}
186
187#[derive(Debug)]
188struct DbInner {
189 columns: Vec<Column>,
190 options: Options,
191 shutdown: AtomicBool,
192 log: Log,
193 commit_queue: Mutex<CommitQueue>,
194 commit_queue_full_cv: Condvar,
195 log_worker_wait: WaitCondvar<bool>,
196 commit_worker_wait: Arc<WaitCondvar<bool>>,
197 commit_overlay: RwLock<Vec<CommitOverlay>>,
199 trees: RwLock<HashMap<ColId, Trees>>,
200 log_queue_wait: WaitCondvar<i64>,
202 flush_worker_wait: Arc<WaitCondvar<bool>>,
203 cleanup_worker_wait: WaitCondvar<bool>,
204 cleanup_queue_wait: WaitCondvar<bool>,
205 iteration_lock: Mutex<()>,
206 last_enacted: AtomicU64,
207 next_reindex: AtomicU64,
208 bg_err: Mutex<Option<Arc<Error>>>,
209 db_version: u32,
210 lock_file: std::fs::File,
211}
212
213#[derive(Debug)]
214struct WaitCondvar<S> {
215 cv: Condvar,
216 work: Mutex<S>,
217}
218
219impl<S: Default> WaitCondvar<S> {
220 fn new() -> Self {
221 WaitCondvar { cv: Condvar::new(), work: Mutex::new(S::default()) }
222 }
223}
224
225impl WaitCondvar<bool> {
226 fn signal(&self) {
227 let mut work = self.work.lock();
228 *work = true;
229 self.cv.notify_one();
230 }
231
232 pub fn wait(&self) {
233 let mut work = self.work.lock();
234 while !*work {
235 self.cv.wait(&mut work)
236 }
237 *work = false;
238 }
239}
240
241impl DbInner {
242 fn open(options: &Options, opening_mode: OpeningMode) -> Result<DbInner> {
243 if opening_mode == OpeningMode::Create {
244 try_io!(std::fs::create_dir_all(&options.path));
245 } else if !options.path.is_dir() {
246 return Err(Error::DatabaseNotFound)
247 }
248
249 let mut lock_path: std::path::PathBuf = options.path.clone();
250 lock_path.push("lock");
251 let lock_file = try_io!(std::fs::OpenOptions::new()
252 .create(true)
253 .read(true)
254 .write(true)
255 .open(lock_path.as_path()));
256 lock_file.try_lock_exclusive().map_err(Error::Locked)?;
257
258 let metadata = options.load_and_validate_metadata(opening_mode == OpeningMode::Create)?;
259 let mut columns = Vec::with_capacity(metadata.columns.len());
260 let mut commit_overlay = Vec::with_capacity(metadata.columns.len());
261 let log = Log::open(options)?;
262 let last_enacted = log.replay_record_id().unwrap_or(2) - 1;
263 for c in 0..metadata.columns.len() {
264 let column = Column::open(c as ColId, options, &metadata)?;
265 commit_overlay.push(CommitOverlay::new());
266 columns.push(column);
267 }
268 log::debug!(target: "parity-db", "Opened db {:?}, metadata={:?}", options, metadata);
269 let mut options = options.clone();
270 if options.salt.is_none() {
271 options.salt = Some(metadata.salt);
272 }
273
274 Ok(DbInner {
275 columns,
276 options,
277 shutdown: AtomicBool::new(false),
278 log,
279 commit_queue: Mutex::new(Default::default()),
280 commit_queue_full_cv: Condvar::new(),
281 log_worker_wait: WaitCondvar::new(),
282 commit_worker_wait: Arc::new(WaitCondvar::new()),
283 commit_overlay: RwLock::new(commit_overlay),
284 trees: RwLock::new(Default::default()),
285 log_queue_wait: WaitCondvar::new(),
286 flush_worker_wait: Arc::new(WaitCondvar::new()),
287 cleanup_worker_wait: WaitCondvar::new(),
288 cleanup_queue_wait: WaitCondvar::new(),
289 iteration_lock: Mutex::new(()),
290 next_reindex: AtomicU64::new(1),
291 last_enacted: AtomicU64::new(last_enacted),
292 bg_err: Mutex::new(None),
293 db_version: metadata.version,
294 lock_file,
295 })
296 }
297
298 fn init_table_data(&mut self) -> Result<()> {
299 for column in &mut self.columns {
300 column.init_table_data()?;
301 }
302 Ok(())
303 }
304
305 fn get(&self, col: ColId, key: &[u8], external_call: bool) -> Result<Option<Value>> {
306 if self.options.columns[col as usize].multitree && external_call {
307 return Err(Error::InvalidConfiguration(
308 "get not supported for multitree columns.".to_string(),
309 ))
310 }
311 match &self.columns[col as usize] {
312 Column::Hash(column) => {
313 let key = column.hash_key(key);
314 let overlay = self.commit_overlay.read();
315 if let Some(v) = overlay.get(col as usize).and_then(|o| o.get(&key)) {
317 return Ok(v.map(|i| i.as_ref().to_vec()))
318 }
319 std::mem::drop(overlay);
320 let log = self.log.overlays();
322 Ok(column.get(&key, log)?.map(|(v, _rc)| v))
323 },
324 Column::Tree(column) => {
325 let overlay = self.commit_overlay.read();
326 if let Some(l) = overlay.get(col as usize).and_then(|o| o.btree_get(key)) {
327 return Ok(l.map(|i| i.as_ref().to_vec()))
328 }
329 std::mem::drop(overlay);
330 let log = self.log.overlays().read();
332 column.with_locked(|btree| BTreeTable::get(key, &*log, btree))
333 },
334 }
335 }
336
337 fn get_size(&self, col: ColId, key: &[u8]) -> Result<Option<u32>> {
338 if self.options.columns[col as usize].multitree {
339 return Err(Error::InvalidConfiguration(
340 "get_size not supported for multitree columns.".to_string(),
341 ))
342 }
343 match &self.columns[col as usize] {
344 Column::Hash(column) => {
345 let key = column.hash_key(key);
346 let overlay = self.commit_overlay.read();
347 if let Some(l) = overlay.get(col as usize).and_then(|o| o.get_size(&key)) {
349 return Ok(l)
350 }
351 let log = self.log.overlays();
353 column.get_size(&key, log)
354 },
355 Column::Tree(column) => {
356 let overlay = self.commit_overlay.read();
357 if let Some(l) = overlay.get(col as usize).and_then(|o| o.btree_get(key)) {
358 return Ok(l.map(|v| v.as_ref().len() as u32))
359 }
360 let log = self.log.overlays().read();
361 let l = column.with_locked(|btree| BTreeTable::get(key, &*log, btree))?;
362 Ok(l.map(|v| v.len() as u32))
363 },
364 }
365 }
366
367 fn get_root(&self, col: ColId, key: &[u8]) -> Result<Option<(Vec<u8>, Children)>> {
368 if !self.options.columns[col as usize].multitree {
369 return Err(Error::InvalidConfiguration("Not a multitree column.".to_string()))
370 }
371 if !self.options.columns[col as usize].append_only &&
372 !self.options.columns[col as usize].allow_direct_node_access
373 {
374 return Err(Error::InvalidConfiguration(
375 "get_root can only be called on a column with append_only or allow_direct_node_access options.".to_string(),
376 ))
377 }
378 let value = self.get(col, key, false)?;
379 if let Some(data) = value {
380 return Ok(Some(unpack_node_data(data)?))
381 }
382 Ok(None)
383 }
384
385 fn get_node(
386 &self,
387 col: ColId,
388 node_address: NodeAddress,
389 external_call: bool,
390 ) -> Result<Option<(Vec<u8>, Children)>> {
391 if !self.options.columns[col as usize].multitree {
392 return Err(Error::InvalidConfiguration("Not a multitree column.".to_string()))
393 }
394 if !self.options.columns[col as usize].append_only &&
395 !self.options.columns[col as usize].allow_direct_node_access &&
396 external_call
397 {
398 return Err(Error::InvalidConfiguration(
399 "get_node can only be called on a column with append_only or allow_direct_node_access options.".to_string(),
400 ))
401 }
402 match &self.columns[col as usize] {
403 Column::Hash(column) => {
404 let overlay = self.commit_overlay.read();
405 if let Some(v) = overlay.get(col as usize).and_then(|o| o.get_address(node_address))
407 {
408 return Ok(Some(unpack_node_data(v.as_ref().to_vec())?))
409 }
410 let log = self.log.overlays();
411 let value = column.get_value(Address::from_u64(node_address), log)?;
412 if let Some(data) = value {
413 return Ok(Some(unpack_node_data(data)?))
414 }
415 Ok(None)
416 },
417 Column::Tree(_) => Err(Error::InvalidConfiguration("Not a HashColumn.".to_string())),
418 }
419 }
420
421 fn get_node_children(
422 &self,
423 col: ColId,
424 node_address: NodeAddress,
425 external_call: bool,
426 ) -> Result<Option<Children>> {
427 if !self.options.columns[col as usize].multitree {
428 return Err(Error::InvalidConfiguration("Not a multitree column.".to_string()))
429 }
430 if !self.options.columns[col as usize].append_only &&
431 !self.options.columns[col as usize].allow_direct_node_access &&
432 external_call
433 {
434 return Err(Error::InvalidConfiguration(
435 "get_node_children can only be called on a column with append_only or allow_direct_node_access options.".to_string(),
436 ))
437 }
438 match &self.columns[col as usize] {
439 Column::Hash(column) => {
440 let overlay = self.commit_overlay.read();
441 if let Some(v) = overlay.get(col as usize).and_then(|o| o.get_address(node_address))
443 {
444 return Ok(Some(unpack_node_children(v.as_ref())?))
445 }
446 let log = self.log.overlays();
447 let value = column.get_value(Address::from_u64(node_address), log)?;
448 if let Some(data) = value {
449 return Ok(Some(unpack_node_children(&data)?))
450 }
451 Ok(None)
452 },
453 Column::Tree(_) => Err(Error::InvalidConfiguration("Not a HashColumn.".to_string())),
454 }
455 }
456
457 fn get_tree(
458 &self,
459 db: &Arc<DbInner>,
460 col: ColId,
461 key: &[u8],
462 check_existence: bool,
463 ) -> Result<Option<Arc<RwLock<Box<dyn TreeReader + Send + Sync>>>>> {
464 if !self.options.columns[col as usize].multitree {
465 return Err(Error::InvalidConfiguration("Not a multitree column.".to_string()))
466 }
467 match &self.columns[col as usize] {
468 Column::Hash(column) => {
469 if check_existence {
472 let root = self.get(col, key, false).unwrap();
473 if root.is_none() {
474 return Ok(None)
475 }
476 }
477
478 let hash_key = column.hash_key(key);
479
480 let trees = self.trees.upgradable_read();
481
482 if let Some(column_trees) = trees.get(&col) {
483 if let Some(reader) = column_trees.readers.get(&hash_key) {
484 let reader = reader.upgrade();
485 if let Some(reader) = reader {
486 return Ok(Some(reader))
487 }
488 }
489 }
490
491 let mut trees = RwLockUpgradableReadGuard::upgrade(trees);
492
493 let column_trees = trees.entry(col).or_insert_with(|| Trees {
494 readers: Default::default(),
495 to_dereference: Default::default(),
496 });
497
498 let reader: Box<dyn TreeReader + Send + Sync> =
499 Box::new(DbTreeReader { db: db.clone(), col, key: hash_key });
500 let reader = Arc::new(RwLock::new(reader));
501
502 column_trees.readers.insert(hash_key, Arc::downgrade(&reader));
503
504 Ok(Some(reader))
505 },
506 Column::Tree(_) => Err(Error::InvalidConfiguration("Not a HashColumn.".to_string())),
507 }
508 }
509
510 fn btree_iter(&self, col: ColId) -> Result<BTreeIterator<'_>> {
511 match &self.columns[col as usize] {
512 Column::Hash(_column) =>
513 Err(Error::InvalidConfiguration("Not an indexed column.".to_string())),
514 Column::Tree(column) => {
515 let log = self.log.overlays();
516 BTreeIterator::new(column, col, log, &self.commit_overlay)
517 },
518 }
519 }
520
521 fn commit<I, K>(&self, tx: I) -> Result<()>
524 where
525 I: IntoIterator<Item = (ColId, K, Option<Value>)>,
526 K: AsRef<[u8]>,
527 {
528 self.commit_changes(tx.into_iter().map(|(c, k, v)| {
529 (
530 c,
531 match v {
532 Some(v) => Operation::Set(k.as_ref().to_vec(), v),
533 None => Operation::Dereference(k.as_ref().to_vec()),
534 },
535 )
536 }))
537 }
538
539 fn commit_changes<I, V>(&self, tx: I) -> Result<()>
540 where
541 I: IntoIterator<Item = (ColId, Operation<Vec<u8>, V>)>,
542 V: Into<RcValue>,
543 {
544 let mut commit: CommitChangeSet = Default::default();
545 for (col, change) in tx.into_iter() {
546 if self.options.columns[col as usize].btree_index {
547 commit
548 .btree_indexed
549 .entry(col)
550 .or_insert_with(|| BTreeChangeSet::new(col))
551 .push(change)?
552 } else if self.options.columns[col as usize].multitree {
553 match &self.columns[col as usize] {
554 Column::Hash(column) =>
555 match change {
556 Operation::Set(..) |
557 Operation::Reference(..) |
558 Operation::Dereference(..) =>
559 return Err(Error::InvalidConfiguration(
560 "Invalid operation for multitree column".to_string(),
561 )),
562 Operation::InsertTree(..) => {
563 let (root_data, node_values) = column.claim_tree_values(&change)?;
564
565 let trees = self.trees.read();
566 if let Some(column_trees) = trees.get(&col) {
567 for (hash, count) in &column_trees.to_dereference {
568 assert!(*count > 0);
569
570 let mut tree_active = false;
572 if let Some(reader) = column_trees.readers.get(hash) {
573 let reader = reader.upgrade();
574 if let Some(reader) = reader {
575 if reader.is_locked() {
576 tree_active = true;
577 }
578 }
579 }
580 if tree_active {
581 commit
582 .indexed
583 .entry(col)
584 .or_insert_with(|| IndexedChangeSet::new(col))
585 .used_trees
586 .insert(*hash);
587 }
588 }
589 }
590 drop(trees);
591
592 let root_operation = Operation::Set(change.key(), root_data);
593 commit
594 .indexed
595 .entry(col)
596 .or_insert_with(|| IndexedChangeSet::new(col))
597 .push(root_operation, &self.options, self.db_version)?;
598
599 for node_change in node_values {
600 commit
601 .indexed
602 .entry(col)
603 .or_insert_with(|| IndexedChangeSet::new(col))
604 .push_node_change(node_change);
605 }
606 },
607 Operation::ReferenceTree(..) => {
608 if !self.options.columns[col as usize].append_only {
609 let root_operation =
610 Operation::<&_, V>::Reference(change.key());
611 commit
612 .indexed
613 .entry(col)
614 .or_insert_with(|| IndexedChangeSet::new(col))
615 .push(root_operation, &self.options, self.db_version)?;
616 }
617 },
618 Operation::DereferenceTree(key) => {
619 if self.options.columns[col as usize].append_only {
620 return Err(Error::InvalidConfiguration("Attempting to dereference a tree from an append_only column.".to_string()))
621 }
622 let value = self.get(col, &key, false)?;
623 if let Some(data) = value {
624 let root_data = unpack_node_data(data)?;
625 let children = root_data.1;
626 let salt = self.options.salt.unwrap_or_default();
627 let hash = hash_key(
628 &key,
629 &salt,
630 self.options.columns[col as usize].uniform,
631 self.db_version,
632 );
633
634 let mut trees = self.trees.write();
635
636 let column_trees = trees.entry(col).or_insert_with(|| Trees {
637 readers: Default::default(),
638 to_dereference: Default::default(),
639 });
640 let count =
641 column_trees.to_dereference.get(&hash).unwrap_or(&0) + 1;
642 column_trees.to_dereference.insert(hash, count);
643
644 drop(trees);
645
646 commit.check_for_deferral = true;
647
648 let node_change =
649 NodeChange::DereferenceChildren(key, hash, children);
650
651 commit
652 .indexed
653 .entry(col)
654 .or_insert_with(|| IndexedChangeSet::new(col))
655 .push_node_change(node_change);
656 } else {
657 return Err(Error::InvalidConfiguration(
658 "No entry for tree root".to_string(),
659 ))
660 }
661 },
662 },
663 Column::Tree(_) =>
664 return Err(Error::InvalidConfiguration("Not a HashColumn".to_string())),
665 }
666 } else {
667 commit.indexed.entry(col).or_insert_with(|| IndexedChangeSet::new(col)).push(
668 change,
669 &self.options,
670 self.db_version,
671 )?
672 }
673 }
674
675 self.commit_raw(commit)
676 }
677
678 fn commit_raw(&self, commit: CommitChangeSet) -> Result<()> {
679 let mut queue = self.commit_queue.lock();
680
681 #[cfg(any(test, feature = "instrumentation"))]
682 let might_wait_because_the_queue_is_full = self.options.with_background_thread;
683 #[cfg(not(any(test, feature = "instrumentation")))]
684 let might_wait_because_the_queue_is_full = true;
685 if might_wait_because_the_queue_is_full && queue.bytes > MAX_COMMIT_QUEUE_BYTES {
686 log::debug!(target: "parity-db", "Waiting, queue size={}", queue.bytes);
687 self.commit_queue_full_cv.wait(&mut queue);
688 }
689
690 {
691 let bg_err = self.bg_err.lock();
692 if let Some(err) = &*bg_err {
693 return Err(Error::Background(err.clone()))
694 }
695 }
696
697 let mut overlay = self.commit_overlay.write();
698
699 queue.record_id += 1;
700 let record_id = queue.record_id;
701
702 let mut bytes = 0;
703 for (c, indexed) in &commit.indexed {
704 indexed.copy_to_overlay(
705 &mut overlay[*c as usize],
706 record_id,
707 &mut bytes,
708 &self.options,
709 )?;
710 }
711
712 for (c, iterset) in &commit.btree_indexed {
713 iterset.copy_to_overlay(
714 &mut overlay[*c as usize].btree_indexed,
715 record_id,
716 &mut bytes,
717 &self.options,
718 )?;
719 }
720
721 let commit = Commit { id: record_id, changeset: commit, bytes };
722
723 log::debug!(
724 target: "parity-db",
725 "Queued commit {}, {} bytes",
726 commit.id,
727 bytes,
728 );
729 queue.commits.push_back(commit);
730 queue.bytes += bytes;
731 self.log_worker_wait.signal();
732 Ok(())
733 }
734
735 fn defer_commit(
736 &self,
737 mut queue: MutexGuard<CommitQueue>,
738 mut commit: CommitChangeSet,
739 old_bytes: usize,
740 old_id: u64,
741 new_id: Option<u64>,
742 ) -> Result<()> {
743 let record_id = if let Some(id) = new_id {
744 id
745 } else {
746 queue.record_id += 1;
747 queue.record_id
748 };
749
750 let bytes = if record_id != old_id {
751 let mut overlay = self.commit_overlay.write();
752
753 let mut bytes = 0;
754
755 for (c, indexed) in &commit.indexed {
756 indexed.copy_to_overlay(
757 &mut overlay[*c as usize],
758 record_id,
759 &mut bytes,
760 &self.options,
761 )?;
762 }
763
764 for (c, iterset) in &commit.btree_indexed {
765 iterset.copy_to_overlay(
766 &mut overlay[*c as usize].btree_indexed,
767 record_id,
768 &mut bytes,
769 &self.options,
770 )?;
771 }
772
773 {
774 for (c, key_values) in commit.indexed.iter() {
776 key_values.clean_overlay(&mut overlay[*c as usize], old_id);
777 }
778 for (c, iterset) in commit.btree_indexed.iter_mut() {
779 iterset.clean_overlay(&mut overlay[*c as usize].btree_indexed, old_id);
780 }
781 }
782
783 bytes
784 } else {
785 old_bytes
786 };
787
788 let commit = Commit { id: record_id, changeset: commit, bytes };
789
790 log::debug!(
791 target: "parity-db",
792 "Deferred commit, old id: {}, new id: {}",
793 old_id,
794 record_id,
795 );
796 queue.commits.push_back(commit);
797 queue.bytes += bytes;
798 Ok(())
799 }
800
801 fn process_commits(&self, db: &Arc<DbInner>) -> Result<bool> {
802 #[cfg(any(test, feature = "instrumentation"))]
803 let might_wait_because_the_queue_is_full = self.options.with_background_thread;
804 #[cfg(not(any(test, feature = "instrumentation")))]
805 let might_wait_because_the_queue_is_full = true;
806 if might_wait_because_the_queue_is_full {
807 let mut queue = self.log_queue_wait.work.lock();
809 if !self.shutdown.load(Ordering::Relaxed) && *queue > MAX_LOG_QUEUE_BYTES {
810 log::debug!(target: "parity-db", "Waiting, log_bytes={}", queue);
811 self.log_queue_wait.cv.wait(&mut queue);
812 }
813 }
814 let commit = {
815 let mut queue = self.commit_queue.lock();
816 if let Some(commit) = queue.commits.pop_front() {
817 queue.bytes -= commit.bytes;
818 log::debug!(
819 target: "parity-db",
820 "Removed {}. Still queued commits {} bytes",
821 commit.bytes,
822 queue.bytes,
823 );
824 if queue.bytes <= MAX_COMMIT_QUEUE_BYTES &&
825 (queue.bytes + commit.bytes) > MAX_COMMIT_QUEUE_BYTES
826 {
827 log::debug!(
829 target: "parity-db",
830 "Waking up commit queue worker",
831 );
832 self.commit_queue_full_cv.notify_all();
833 }
834 Some(commit)
835 } else {
836 None
837 }
838 };
839
840 if let Some(mut commit) = commit {
841 if commit.changeset.check_for_deferral {
842 let mut defer = false;
843 'outer: for (col, key_values) in commit.changeset.indexed.iter() {
844 for change in &key_values.node_changes {
845 if let NodeChange::DereferenceChildren(_key, hash, _children) = change {
846 let trees = self.trees.read();
849 if let Some(column_trees) = trees.get(&col) {
850 let mut tree_active = false;
851 if let Some(reader) = column_trees.readers.get(hash) {
852 let reader = reader.upgrade();
853 if let Some(reader) = reader {
854 if reader.is_locked() {
855 tree_active = true;
856 }
857 }
858 }
859 if tree_active {
860 defer = true;
861 break 'outer
862 }
863 }
864 drop(trees);
865
866 let queue = self.commit_queue.lock();
869 for commit in &queue.commits {
870 for (_col, change_set) in &commit.changeset.indexed {
871 for tree in &change_set.used_trees {
872 if tree == hash {
873 defer = true;
874 break 'outer
875 }
876 }
877 }
878 }
879 }
880 }
881 }
882 if defer {
883 let queue = self.commit_queue.lock();
884 let new_id = if queue.commits.len() > 0 {
885 None
887 } else {
888 Some(commit.id)
890 };
891 self.defer_commit(queue, commit.changeset, commit.bytes, commit.id, new_id)?;
892
893 return Ok(true)
894 } else {
895 for (col, key_values) in commit.changeset.indexed.iter() {
896 for change in &key_values.node_changes {
897 if let NodeChange::DereferenceChildren(_key, hash, _children) = change {
898 let mut trees = self.trees.write();
899 if let Some(column_trees) = trees.get_mut(&col) {
900 let count = column_trees.to_dereference.get(hash).unwrap_or(&0);
901 assert!(*count > 0);
902 if *count == 1 {
903 column_trees.to_dereference.remove(hash);
904 } else {
905 column_trees.to_dereference.insert(*hash, count - 1);
906 }
907 }
908 }
909 }
910 }
911 }
912 }
913
914 let mut reindex = false;
915 let mut writer = self.log.begin_record();
916 log::debug!(
917 target: "parity-db",
918 "Processing commit {}, record {}, {} bytes",
919 commit.id,
920 writer.record_id(),
921 commit.bytes,
922 );
923 let mut ops: u64 = 0;
924 for (c, key_values) in commit.changeset.indexed.iter() {
925 key_values.write_plan(
926 db,
927 *c,
928 &self.columns[*c as usize],
929 &mut writer,
930 &mut ops,
931 &mut reindex,
932 )?;
933 }
934
935 for (c, btree) in commit.changeset.btree_indexed.iter_mut() {
936 match &self.columns[*c as usize] {
937 Column::Hash(_column) =>
938 return Err(Error::InvalidConfiguration(
939 "Not an indexed column.".to_string(),
940 )),
941 Column::Tree(column) => {
942 btree.write_plan(column, &mut writer, &mut ops)?;
943 },
944 }
945 }
946
947 for c in self.columns.iter() {
949 c.complete_plan(&mut writer)?;
950 }
951 let record_id = writer.record_id();
952 let l = writer.drain();
953
954 let bytes = {
955 let bytes = self.log.end_record(l)?;
956 let mut logged_bytes = self.log_queue_wait.work.lock();
957 *logged_bytes += bytes as i64;
958 self.flush_worker_wait.signal();
959 bytes
960 };
961
962 {
963 let mut overlay = self.commit_overlay.write();
965 for (c, key_values) in commit.changeset.indexed.iter() {
966 key_values.clean_overlay(&mut overlay[*c as usize], commit.id);
967 }
968 for (c, iterset) in commit.changeset.btree_indexed.iter_mut() {
969 iterset.clean_overlay(&mut overlay[*c as usize].btree_indexed, commit.id);
970 }
971 }
972
973 if reindex {
974 self.start_reindex(record_id);
975 }
976
977 log::debug!(
978 target: "parity-db",
979 "Processed commit {} (record {}), {} ops, {} bytes written",
980 commit.id,
981 record_id,
982 ops,
983 bytes,
984 );
985 Ok(true)
986 } else {
987 Ok(false)
988 }
989 }
990
991 fn start_reindex(&self, record_id: u64) {
992 log::trace!(target: "parity-db", "Scheduled reindex at record {}", record_id);
993 self.next_reindex.store(record_id, Ordering::SeqCst);
994 }
995
996 fn process_reindex(&self) -> Result<bool> {
997 let next_reindex = self.next_reindex.load(Ordering::SeqCst);
998 if next_reindex == 0 || next_reindex > self.last_enacted.load(Ordering::SeqCst) {
999 return Ok(false)
1000 }
1001 for column in self.columns.iter() {
1003 let column = if let Column::Hash(c) = column { c } else { continue };
1004 let ReindexBatch {
1005 drop_index,
1006 batch,
1007 drop_ref_count,
1008 ref_count_batch,
1009 ref_count_batch_source,
1010 } = column.reindex(&self.log)?;
1011 if !batch.is_empty() || drop_index.is_some() {
1012 debug_assert!(
1013 ref_count_batch.is_empty() &&
1014 ref_count_batch_source.is_none() &&
1015 drop_ref_count.is_none()
1016 );
1017 let mut next_reindex = false;
1018 let mut writer = self.log.begin_record();
1019 log::debug!(
1020 target: "parity-db",
1021 "Creating reindex record {}",
1022 writer.record_id(),
1023 );
1024 for (key, address) in batch.into_iter() {
1025 if let PlanOutcome::NeedReindex =
1026 column.write_reindex_plan(&key, address, &mut writer)?
1027 {
1028 next_reindex = true
1029 }
1030 }
1031 if let Some(table) = drop_index {
1032 writer.drop_table(table);
1033 }
1034 let record_id = writer.record_id();
1035 let l = writer.drain();
1036
1037 let mut logged_bytes = self.log_queue_wait.work.lock();
1038 let bytes = self.log.end_record(l)?;
1039 log::debug!(
1040 target: "parity-db",
1041 "Created reindex record {}, {} bytes",
1042 record_id,
1043 bytes,
1044 );
1045 *logged_bytes += bytes as i64;
1046 if next_reindex {
1047 self.start_reindex(record_id);
1048 }
1049 self.flush_worker_wait.signal();
1050 return Ok(true)
1051 }
1052 if !ref_count_batch.is_empty() || drop_ref_count.is_some() {
1053 debug_assert!(batch.is_empty() && drop_index.is_none());
1054 debug_assert!(ref_count_batch_source.is_some());
1055 let ref_count_source = ref_count_batch_source.unwrap();
1056 let mut next_reindex = false;
1057 let mut writer = self.log.begin_record();
1058 log::debug!(
1059 target: "parity-db",
1060 "Creating ref count reindex record {}",
1061 writer.record_id(),
1062 );
1063 for (address, ref_count) in ref_count_batch.into_iter() {
1064 if let PlanOutcome::NeedReindex = column.write_ref_count_reindex_plan(
1065 address,
1066 ref_count,
1067 ref_count_source,
1068 &mut writer,
1069 )? {
1070 next_reindex = true
1071 }
1072 }
1073 if let Some(table) = drop_ref_count {
1074 writer.drop_ref_count_table(table);
1075 }
1076 let record_id = writer.record_id();
1077 let l = writer.drain();
1078
1079 let mut logged_bytes = self.log_queue_wait.work.lock();
1080 let bytes = self.log.end_record(l)?;
1081 log::debug!(
1082 target: "parity-db",
1083 "Created ref count reindex record {}, {} bytes",
1084 record_id,
1085 bytes,
1086 );
1087 *logged_bytes += bytes as i64;
1088 if next_reindex {
1089 self.start_reindex(record_id);
1090 }
1091 self.flush_worker_wait.signal();
1092 return Ok(true)
1093 }
1094 }
1095 self.next_reindex.store(0, Ordering::SeqCst);
1096 Ok(false)
1097 }
1098
1099 fn enact_logs(&self, validation_mode: bool) -> Result<bool> {
1100 let _iteration_lock = self.iteration_lock.lock();
1101 let cleared = {
1102 let reader = match self.log.read_next(validation_mode) {
1103 Ok(reader) => reader,
1104 Err(Error::Corruption(_)) if validation_mode => {
1105 log::debug!(target: "parity-db", "Bad log header");
1106 self.log.clear_replay_logs();
1107 return Ok(false)
1108 },
1109 Err(e) => return Err(e),
1110 };
1111 if let Some(mut reader) = reader {
1112 log::debug!(
1113 target: "parity-db",
1114 "Enacting log record {}",
1115 reader.record_id(),
1116 );
1117 if validation_mode {
1118 if reader.record_id() != self.last_enacted.load(Ordering::Relaxed) + 1 {
1119 log::warn!(
1120 target: "parity-db",
1121 "Log sequence error. Expected record {}, got {}",
1122 self.last_enacted.load(Ordering::Relaxed) + 1,
1123 reader.record_id(),
1124 );
1125 drop(reader);
1126 self.log.clear_replay_logs();
1127 return Ok(false)
1128 }
1129 loop {
1131 let next = match reader.next() {
1132 Ok(next) => next,
1133 Err(e) => {
1134 log::debug!(target: "parity-db", "Error reading log: {:?}", e);
1135 return Ok(false)
1136 },
1137 };
1138 match next {
1139 LogAction::BeginRecord => {
1140 log::debug!(target: "parity-db", "Unexpected log header");
1141 drop(reader);
1142 self.log.clear_replay_logs();
1143 return Ok(false)
1144 },
1145 LogAction::EndRecord => break,
1146 LogAction::InsertIndex(insertion) => {
1147 let col = insertion.table.col() as usize;
1148 if let Err(e) = self.columns.get(col).map_or_else(
1149 || Err(Error::Corruption(format!("Invalid column id {col}"))),
1150 |col| {
1151 col.validate_plan(
1152 LogAction::InsertIndex(insertion),
1153 &mut reader,
1154 )
1155 },
1156 ) {
1157 log::warn!(target: "parity-db", "Error validating log: {:?}.", e);
1158 drop(reader);
1159 self.log.clear_replay_logs();
1160 return Ok(false)
1161 }
1162 },
1163 LogAction::InsertValue(insertion) => {
1164 let col = insertion.table.col() as usize;
1165 if let Err(e) = self.columns.get(col).map_or_else(
1166 || Err(Error::Corruption(format!("Invalid column id {col}"))),
1167 |col| {
1168 col.validate_plan(
1169 LogAction::InsertValue(insertion),
1170 &mut reader,
1171 )
1172 },
1173 ) {
1174 log::warn!(target: "parity-db", "Error validating log: {:?}.", e);
1175 drop(reader);
1176 self.log.clear_replay_logs();
1177 return Ok(false)
1178 }
1179 },
1180 LogAction::InsertRefCount(insertion) => {
1181 let col = insertion.table.col() as usize;
1182 if let Err(e) = self.columns.get(col).map_or_else(
1183 || Err(Error::Corruption(format!("Invalid column id {col}"))),
1184 |col| {
1185 col.validate_plan(
1186 LogAction::InsertRefCount(insertion),
1187 &mut reader,
1188 )
1189 },
1190 ) {
1191 log::warn!(target: "parity-db", "Error validating log: {:?}.", e);
1192 drop(reader);
1193 self.log.clear_replay_logs();
1194 return Ok(false)
1195 }
1196 },
1197 LogAction::DropTable(_) | LogAction::DropRefCountTable(_) => continue,
1198 }
1199 }
1200 reader.reset()?;
1201 reader.next()?;
1202 }
1203 loop {
1204 match reader.next()? {
1205 LogAction::BeginRecord =>
1206 return Err(Error::Corruption("Bad log record".into())),
1207 LogAction::EndRecord => break,
1208 LogAction::InsertIndex(insertion) => {
1209 self.columns[insertion.table.col() as usize]
1210 .enact_plan(LogAction::InsertIndex(insertion), &mut reader)?;
1211 },
1212 LogAction::InsertValue(insertion) => {
1213 self.columns[insertion.table.col() as usize]
1214 .enact_plan(LogAction::InsertValue(insertion), &mut reader)?;
1215 },
1216 LogAction::InsertRefCount(insertion) => {
1217 self.columns[insertion.table.col() as usize]
1218 .enact_plan(LogAction::InsertRefCount(insertion), &mut reader)?;
1219 },
1220 LogAction::DropTable(id) => {
1221 log::debug!(
1222 target: "parity-db",
1223 "Dropping index {}",
1224 id,
1225 );
1226 match &self.columns[id.col() as usize] {
1227 Column::Hash(col) => {
1228 col.drop_index(id)?;
1229 self.start_reindex(reader.record_id());
1231 },
1232 Column::Tree(_) => (),
1233 }
1234 },
1235 LogAction::DropRefCountTable(id) => {
1236 log::debug!(
1237 target: "parity-db",
1238 "Dropping ref count {}",
1239 id,
1240 );
1241 match &self.columns[id.col() as usize] {
1242 Column::Hash(col) => {
1243 col.drop_ref_count(id)?;
1244 self.start_reindex(reader.record_id());
1246 },
1247 Column::Tree(_) => (),
1248 }
1249 },
1250 }
1251 }
1252 log::debug!(
1253 target: "parity-db",
1254 "Enacted log record {}, {} bytes",
1255 reader.record_id(),
1256 reader.read_bytes(),
1257 );
1258 let record_id = reader.record_id();
1259 let bytes = reader.read_bytes();
1260 let cleared = reader.drain();
1261 self.last_enacted.store(record_id, Ordering::SeqCst);
1262 Some((record_id, cleared, bytes))
1263 } else {
1264 log::debug!(target: "parity-db", "End of log");
1265 None
1266 }
1267 };
1268
1269 if let Some((record_id, cleared, bytes)) = cleared {
1270 self.log.end_read(cleared, record_id);
1271 {
1272 if !validation_mode {
1273 let mut queue = self.log_queue_wait.work.lock();
1274 if *queue < bytes as i64 {
1275 log::warn!(
1276 target: "parity-db",
1277 "Detected log underflow record {}, {} bytes, {} queued, reindex = {}",
1278 record_id,
1279 bytes,
1280 *queue,
1281 self.next_reindex.load(Ordering::SeqCst),
1282 );
1283 }
1284 *queue -= bytes as i64;
1285 if *queue <= MAX_LOG_QUEUE_BYTES &&
1286 (*queue + bytes as i64) > MAX_LOG_QUEUE_BYTES
1287 {
1288 self.log_queue_wait.cv.notify_one();
1289 }
1290 log::debug!(target: "parity-db", "Log queue size: {} bytes", *queue);
1291 }
1292
1293 let max_logs = if self.options.sync_data { MAX_LOG_FILES } else { KEEP_LOGS };
1294 let dirty_logs = self.log.num_dirty_logs();
1295 if !validation_mode {
1296 while !self.shutdown.load(Ordering::Relaxed) &&
1297 self.log.num_dirty_logs() > max_logs
1298 {
1299 log::debug!(target: "parity-db", "Waiting for log cleanup. Queued: {}", dirty_logs);
1300 self.cleanup_worker_wait.signal();
1301 self.cleanup_queue_wait.wait();
1302 }
1303 }
1304 }
1305 Ok(true)
1306 } else {
1307 Ok(false)
1308 }
1309 }
1310
1311 fn flush_logs(&self, min_log_size: u64) -> Result<bool> {
1312 let has_flushed = self.log.flush_one(min_log_size)?;
1313 if has_flushed {
1314 self.commit_worker_wait.signal();
1315 }
1316 Ok(has_flushed)
1317 }
1318
1319 fn clean_logs(&self) -> Result<bool> {
1320 let keep_logs = if self.options.sync_data { 0 } else { KEEP_LOGS };
1321 let num_cleanup = self.log.num_dirty_logs();
1322 let result = if num_cleanup > keep_logs {
1323 if self.options.sync_data {
1324 for c in self.columns.iter() {
1325 c.flush()?;
1326 }
1327 }
1328 self.log.clean_logs(num_cleanup - keep_logs)?
1329 } else {
1330 false
1331 };
1332 self.cleanup_queue_wait.signal();
1333 Ok(result)
1334 }
1335
1336 fn clean_all_logs(&self) -> Result<()> {
1337 for c in self.columns.iter() {
1338 c.flush()?;
1339 }
1340 let num_cleanup = self.log.num_dirty_logs();
1341 self.log.clean_logs(num_cleanup)?;
1342 Ok(())
1343 }
1344
1345 fn replay_all_logs(&self) -> Result<()> {
1346 while let Some(id) = self.log.replay_next()? {
1347 log::debug!(target: "parity-db", "Replaying database log {}", id);
1348 while self.enact_logs(true)? {}
1349 }
1350
1351 for c in self.columns.iter() {
1353 c.refresh_metadata()?;
1354 }
1355 log::debug!(target: "parity-db", "Replay is complete.");
1356 Ok(())
1357 }
1358
1359 fn shutdown(&self) {
1360 self.shutdown.store(true, Ordering::SeqCst);
1361 self.log_queue_wait.cv.notify_one();
1362 self.flush_worker_wait.signal();
1363 self.log_worker_wait.signal();
1364 self.commit_worker_wait.signal();
1365 self.cleanup_worker_wait.signal();
1366 }
1367
1368 fn kill_logs(&self, db: &Arc<DbInner>) -> Result<()> {
1369 {
1370 if let Some(err) = self.bg_err.lock().as_ref() {
1371 log::debug!(target: "parity-db", "Shutdown with error state {}", err);
1374 self.log.clean_logs(self.log.num_dirty_logs())?;
1375 return Ok(())
1376 }
1377 }
1378 log::debug!(target: "parity-db", "Processing leftover commits");
1379 while self.enact_logs(false)? {}
1381 self.flush_logs(0)?;
1382 while self.process_commits(db)? {}
1383 while self.enact_logs(false)? {}
1384 self.flush_logs(0)?;
1385 while self.enact_logs(false)? {}
1386 self.clean_all_logs()?;
1387 self.log.kill_logs()?;
1388 if self.options.stats {
1389 let mut path = self.options.path.clone();
1390 path.push("stats.txt");
1391 match std::fs::File::create(path) {
1392 Ok(file) => {
1393 let mut writer = std::io::BufWriter::new(file);
1394 if let Err(e) = self.write_stats_text(&mut writer, None) {
1395 log::warn!(target: "parity-db", "Error writing stats file: {:?}", e)
1396 }
1397 },
1398 Err(e) => log::warn!(target: "parity-db", "Error creating stats file: {:?}", e),
1399 }
1400 }
1401 Ok(())
1402 }
1403
1404 fn write_stats_text(&self, writer: &mut impl std::io::Write, column: Option<u8>) -> Result<()> {
1405 if let Some(col) = column {
1406 self.columns[col as usize].write_stats_text(writer)
1407 } else {
1408 for c in self.columns.iter() {
1409 c.write_stats_text(writer)?;
1410 }
1411 Ok(())
1412 }
1413 }
1414
1415 fn clear_stats(&self, column: Option<u8>) -> Result<()> {
1416 if let Some(col) = column {
1417 self.columns[col as usize].clear_stats()
1418 } else {
1419 for c in self.columns.iter() {
1420 c.clear_stats()?;
1421 }
1422 Ok(())
1423 }
1424 }
1425
1426 fn stats(&self) -> StatSummary {
1427 StatSummary { columns: self.columns.iter().map(|c| c.stats()).collect() }
1428 }
1429
1430 fn store_err(&self, result: Result<()>) {
1431 if let Err(e) = result {
1432 log::warn!(target: "parity-db", "Background worker error: {}", e);
1433 let mut err = self.bg_err.lock();
1434 if err.is_none() {
1435 *err = Some(Arc::new(e));
1436 self.shutdown();
1437 }
1438 self.commit_queue_full_cv.notify_all();
1439 }
1440 }
1441
1442 fn iter_column_while(&self, c: ColId, f: impl FnMut(ValueIterState) -> bool) -> Result<()> {
1443 let _lock = self.iteration_lock.lock();
1444 match &self.columns[c as usize] {
1445 Column::Hash(column) => column.iter_values(&self.log, f),
1446 Column::Tree(_) => unimplemented!(),
1447 }
1448 }
1449
1450 fn iter_column_index_while(&self, c: ColId, f: impl FnMut(IterState) -> bool) -> Result<()> {
1451 let _lock = self.iteration_lock.lock();
1452 match &self.columns[c as usize] {
1453 Column::Hash(column) => column.iter_index(&self.log, f),
1454 Column::Tree(_) => unimplemented!(),
1455 }
1456 }
1457}
1458
1459pub struct Db {
1461 inner: Arc<DbInner>,
1462 commit_thread: Option<thread::JoinHandle<()>>,
1463 flush_thread: Option<thread::JoinHandle<()>>,
1464 log_thread: Option<thread::JoinHandle<()>>,
1465 cleanup_thread: Option<thread::JoinHandle<()>>,
1466}
1467
1468impl Db {
1469 #[cfg(test)]
1470 pub(crate) fn with_columns(path: &std::path::Path, num_columns: u8) -> Result<Db> {
1471 let options = Options::with_columns(path, num_columns);
1472 Self::open_inner(&options, OpeningMode::Create)
1473 }
1474
1475 pub fn open(options: &Options) -> Result<Db> {
1478 Self::open_inner(options, OpeningMode::Write)
1479 }
1480
1481 pub fn open_or_create(options: &Options) -> Result<Db> {
1484 Self::open_inner(options, OpeningMode::Create)
1485 }
1486
1487 pub fn open_read_only(options: &Options) -> Result<Db> {
1489 Self::open_inner(options, OpeningMode::ReadOnly)
1490 }
1491
1492 fn open_inner(options: &Options, opening_mode: OpeningMode) -> Result<Db> {
1493 assert!(options.is_valid());
1494 let mut db = DbInner::open(options, opening_mode)?;
1495 if let Err(e) = db.replay_all_logs() {
1498 log::debug!(target: "parity-db", "Error during log replay.");
1499 return Err(e)
1500 } else {
1501 db.log.clear_replay_logs();
1502 db.clean_all_logs()?;
1503 db.log.kill_logs()?;
1504 }
1505 db.init_table_data()?;
1506 let db = Arc::new(db);
1507 #[cfg(any(test, feature = "instrumentation"))]
1508 let start_threads = opening_mode != OpeningMode::ReadOnly && options.with_background_thread;
1509 #[cfg(not(any(test, feature = "instrumentation")))]
1510 let start_threads = opening_mode != OpeningMode::ReadOnly;
1511 let commit_thread = if start_threads {
1512 let commit_worker_db = db.clone();
1513 Some(thread::spawn(move || {
1514 commit_worker_db.store_err(Self::commit_worker(commit_worker_db.clone()))
1515 }))
1516 } else {
1517 None
1518 };
1519 let flush_thread = if start_threads {
1520 let flush_worker_db = db.clone();
1521 #[cfg(any(test, feature = "instrumentation"))]
1522 let min_log_size = if options.always_flush { 0 } else { MIN_LOG_SIZE_BYTES };
1523 #[cfg(not(any(test, feature = "instrumentation")))]
1524 let min_log_size = MIN_LOG_SIZE_BYTES;
1525 Some(thread::spawn(move || {
1526 flush_worker_db.store_err(Self::flush_worker(flush_worker_db.clone(), min_log_size))
1527 }))
1528 } else {
1529 None
1530 };
1531 let log_thread = if start_threads {
1532 let log_worker_db = db.clone();
1533 Some(thread::spawn(move || {
1534 log_worker_db.store_err(Self::log_worker(log_worker_db.clone()))
1535 }))
1536 } else {
1537 None
1538 };
1539 let cleanup_thread = if start_threads {
1540 let cleanup_worker_db = db.clone();
1541 Some(thread::spawn(move || {
1542 cleanup_worker_db.store_err(Self::cleanup_worker(cleanup_worker_db.clone()))
1543 }))
1544 } else {
1545 None
1546 };
1547 Ok(Db { inner: db, commit_thread, flush_thread, log_thread, cleanup_thread })
1548 }
1549
1550 pub fn get(&self, col: ColId, key: &[u8]) -> Result<Option<Value>> {
1552 self.inner.get(col, key, true)
1553 }
1554
1555 pub fn get_size(&self, col: ColId, key: &[u8]) -> Result<Option<u32>> {
1557 self.inner.get_size(col, key)
1558 }
1559
1560 pub fn iter(&self, col: ColId) -> Result<BTreeIterator<'_>> {
1563 self.inner.btree_iter(col)
1564 }
1565
1566 pub fn get_tree(
1567 &self,
1568 col: ColId,
1569 key: &[u8],
1570 ) -> Result<Option<Arc<RwLock<Box<dyn TreeReader + Send + Sync>>>>> {
1571 self.inner.get_tree(&self.inner, col, key, true)
1572 }
1573
1574 pub fn get_root(&self, col: ColId, key: &[u8]) -> Result<Option<(Vec<u8>, Children)>> {
1575 self.inner.get_root(col, key)
1576 }
1577
1578 pub fn get_node(
1579 &self,
1580 col: ColId,
1581 node_address: NodeAddress,
1582 ) -> Result<Option<(Vec<u8>, Children)>> {
1583 self.inner.get_node(col, node_address, true)
1584 }
1585
1586 pub fn get_node_children(
1587 &self,
1588 col: ColId,
1589 node_address: NodeAddress,
1590 ) -> Result<Option<Children>> {
1591 self.inner.get_node_children(col, node_address, true)
1592 }
1593
1594 pub fn commit<I, K>(&self, tx: I) -> Result<()>
1596 where
1597 I: IntoIterator<Item = (ColId, K, Option<Value>)>,
1598 K: AsRef<[u8]>,
1599 {
1600 self.inner.commit(tx)
1601 }
1602
1603 pub fn commit_changes<I>(&self, tx: I) -> Result<()>
1605 where
1606 I: IntoIterator<Item = (ColId, Operation<Vec<u8>, Vec<u8>>)>,
1607 {
1608 self.inner.commit_changes(tx)
1609 }
1610
1611 #[cfg(feature = "arc")]
1615 #[deprecated(
1616 note = "This method will be removed in future versions. Use `commit_changes_bytes` instead"
1617 )]
1618 pub fn commit_changes_shared<I>(&self, tx: I) -> Result<()>
1619 where
1620 I: IntoIterator<Item = (ColId, Operation<Vec<u8>, Arc<Vec<u8>>>)>,
1621 {
1622 self.inner.commit_changes(tx)
1623 }
1624
1625 #[cfg(feature = "bytes")]
1629 pub fn commit_changes_bytes<I>(&self, tx: I) -> Result<()>
1630 where
1631 I: IntoIterator<Item = (ColId, Operation<Vec<u8>, Bytes>)>,
1632 {
1633 self.inner.commit_changes(tx)
1634 }
1635
1636 pub(crate) fn commit_raw(&self, commit: CommitChangeSet) -> Result<()> {
1637 self.inner.commit_raw(commit)
1638 }
1639
1640 pub fn num_columns(&self) -> u8 {
1642 self.inner.columns.len() as u8
1643 }
1644
1645 pub fn iter_column_while(&self, c: ColId, f: impl FnMut(ValueIterState) -> bool) -> Result<()> {
1649 self.inner.iter_column_while(c, f)
1650 }
1651
1652 pub(crate) fn iter_column_index_while(
1657 &self,
1658 c: ColId,
1659 f: impl FnMut(IterState) -> bool,
1660 ) -> Result<()> {
1661 self.inner.iter_column_index_while(c, f)
1662 }
1663
1664 fn commit_worker(db: Arc<DbInner>) -> Result<()> {
1665 let mut more_work = false;
1666 while !db.shutdown.load(Ordering::SeqCst) || more_work {
1667 if !more_work {
1668 db.cleanup_worker_wait.signal();
1669 if !db.log.has_log_files_to_read() {
1670 db.commit_worker_wait.wait();
1671 }
1672 }
1673
1674 more_work = db.enact_logs(false)?;
1675 }
1676 log::debug!(target: "parity-db", "Commit worker shutdown");
1677 Ok(())
1678 }
1679
1680 fn log_worker(db: Arc<DbInner>) -> Result<()> {
1681 let mut more_reindex = db.process_reindex()?;
1683 let mut more_commits = false;
1684 while !db.shutdown.load(Ordering::SeqCst) || more_commits {
1686 if !more_commits && !more_reindex {
1687 db.log_worker_wait.wait();
1688 }
1689
1690 more_commits = db.process_commits(&db)?;
1691 more_reindex = db.process_reindex()?;
1692 }
1693 log::debug!(target: "parity-db", "Log worker shutdown");
1694 Ok(())
1695 }
1696
1697 fn flush_worker(db: Arc<DbInner>, min_log_size: u64) -> Result<()> {
1698 let mut more_work = false;
1699 while !db.shutdown.load(Ordering::SeqCst) {
1700 if !more_work {
1701 db.flush_worker_wait.wait();
1702 }
1703 more_work = db.flush_logs(min_log_size)?;
1704 }
1705 log::debug!(target: "parity-db", "Flush worker shutdown");
1706 Ok(())
1707 }
1708
1709 fn cleanup_worker(db: Arc<DbInner>) -> Result<()> {
1710 let mut more_work = true;
1711 while !db.shutdown.load(Ordering::SeqCst) || more_work {
1712 if !more_work {
1713 db.cleanup_worker_wait.wait();
1714 }
1715 more_work = db.clean_logs()?;
1716 }
1717 log::debug!(target: "parity-db", "Cleanup worker shutdown");
1718 Ok(())
1719 }
1720
1721 pub fn write_stats_text(
1723 &self,
1724 writer: &mut impl std::io::Write,
1725 column: Option<u8>,
1726 ) -> Result<()> {
1727 self.inner.write_stats_text(writer, column)
1728 }
1729
1730 pub fn clear_stats(&self, column: Option<u8>) -> Result<()> {
1732 self.inner.clear_stats(column)
1733 }
1734
1735 pub fn dump(&self, check_param: check::CheckOptions) -> Result<()> {
1737 if let Some(col) = check_param.column {
1738 self.inner.columns[col as usize].dump(&self.inner.log, &check_param, col)?;
1739 } else {
1740 for (ix, c) in self.inner.columns.iter().enumerate() {
1741 c.dump(&self.inner.log, &check_param, ix as ColId)?;
1742 }
1743 }
1744 Ok(())
1745 }
1746
1747 pub fn stats(&self) -> StatSummary {
1749 self.inner.stats()
1750 }
1751
1752 pub fn get_num_column_value_entries(&self, col: ColId) -> Result<u64> {
1753 let column = &self.inner.columns[col as usize];
1754 match column {
1755 Column::Hash(column) => return column.get_num_value_entries(),
1756 Column::Tree(..) =>
1757 return Err(Error::InvalidConfiguration(
1758 "get_num_column_value_entries not implemented for tree columns.".to_string(),
1759 )),
1760 }
1761 }
1762
1763 fn precheck_column_operation(options: &mut Options) -> Result<[u8; 32]> {
1766 let db = Db::open(options)?;
1767 let salt = db.inner.options.salt;
1768 drop(db);
1769 Ok(salt.expect("`salt` is always `Some` after opening the DB; qed"))
1770 }
1771
1772 pub fn add_column(options: &mut Options, new_column_options: ColumnOptions) -> Result<()> {
1774 let salt = Self::precheck_column_operation(options)?;
1775
1776 options.columns.push(new_column_options);
1777 options.write_metadata_with_version(&options.path, &salt, Some(CURRENT_VERSION))?;
1778
1779 Ok(())
1780 }
1781
1782 pub fn drop_last_column(options: &mut Options) -> Result<()> {
1785 let salt = Self::precheck_column_operation(options)?;
1786 let nb_column = options.columns.len();
1787 if nb_column == 0 {
1788 return Ok(())
1789 }
1790 let index = options.columns.len() - 1;
1791 Self::remove_column_files(options, index as u8)?;
1792 options.columns.pop();
1793 options.write_metadata(&options.path, &salt)?;
1794 Ok(())
1795 }
1796
1797 pub fn reset_column(
1800 options: &mut Options,
1801 index: u8,
1802 new_options: Option<ColumnOptions>,
1803 ) -> Result<()> {
1804 let salt = Self::precheck_column_operation(options)?;
1805 Self::remove_column_files(options, index)?;
1806
1807 if let Some(new_options) = new_options {
1808 options.columns[index as usize] = new_options;
1809 options.write_metadata(&options.path, &salt)?;
1810 }
1811
1812 Ok(())
1813 }
1814
1815 fn remove_column_files(options: &mut Options, index: u8) -> Result<()> {
1816 if index as usize >= options.columns.len() {
1817 return Err(Error::IncompatibleColumnConfig {
1818 id: index,
1819 reason: "Column not found".to_string(),
1820 })
1821 }
1822
1823 Column::drop_files(index, options.path.clone())?;
1824 Ok(())
1825 }
1826
1827 #[cfg(feature = "instrumentation")]
1828 pub fn process_reindex(&self) -> Result<()> {
1829 self.inner.process_reindex()?;
1830 Ok(())
1831 }
1832
1833 #[cfg(feature = "instrumentation")]
1834 pub fn process_commits(&self) -> Result<()> {
1835 self.inner.process_commits(&self.inner)?;
1836 Ok(())
1837 }
1838
1839 #[cfg(feature = "instrumentation")]
1840 pub fn flush_logs(&self) -> Result<()> {
1841 self.inner.flush_logs(0)?;
1842 Ok(())
1843 }
1844
1845 #[cfg(feature = "instrumentation")]
1846 pub fn enact_logs(&self) -> Result<()> {
1847 while self.inner.enact_logs(false)? {}
1848 Ok(())
1849 }
1850
1851 #[cfg(feature = "instrumentation")]
1852 pub fn clean_logs(&self) -> Result<()> {
1853 self.inner.clean_logs()?;
1854 Ok(())
1855 }
1856}
1857
1858impl Drop for Db {
1859 fn drop(&mut self) {
1860 self.drop_inner()
1861 }
1862}
1863
1864impl Db {
1865 fn drop_inner(&mut self) {
1866 self.inner.shutdown();
1867 if let Some(t) = self.log_thread.take() {
1868 if let Err(e) = t.join() {
1869 log::warn!(target: "parity-db", "Log thread shutdown error: {:?}", e);
1870 }
1871 }
1872 if let Some(t) = self.flush_thread.take() {
1873 if let Err(e) = t.join() {
1874 log::warn!(target: "parity-db", "Flush thread shutdown error: {:?}", e);
1875 }
1876 }
1877 if let Some(t) = self.commit_thread.take() {
1878 if let Err(e) = t.join() {
1879 log::warn!(target: "parity-db", "Commit thread shutdown error: {:?}", e);
1880 }
1881 }
1882 if let Some(t) = self.cleanup_thread.take() {
1883 if let Err(e) = t.join() {
1884 log::warn!(target: "parity-db", "Cleanup thread shutdown error: {:?}", e);
1885 }
1886 }
1887 if let Err(e) = self.inner.kill_logs(&self.inner) {
1888 log::warn!(target: "parity-db", "Shutdown error: {:?}", e);
1889 }
1890 if let Err(e) = fs2::FileExt::unlock(&self.inner.lock_file) {
1891 log::debug!(target: "parity-db", "Error removing file lock: {:?}", e);
1892 }
1893 }
1894}
1895
1896pub trait TreeReader {
1899 fn get_root(&self) -> Result<Option<(Vec<u8>, Children)>>;
1900 fn get_node(&self, node_address: NodeAddress) -> Result<Option<(Vec<u8>, Children)>>;
1901 fn get_node_children(&self, node_address: NodeAddress) -> Result<Option<Children>>;
1902}
1903
1904#[derive(Debug)]
1905pub struct DbTreeReader {
1906 db: Arc<DbInner>,
1907 col: ColId,
1908 key: Key,
1909}
1910
1911impl TreeReader for DbTreeReader {
1912 fn get_root(&self) -> Result<Option<(Vec<u8>, Children)>> {
1913 match &self.db.columns[self.col as usize] {
1920 Column::Hash(column) => {
1921 let overlay = self.db.commit_overlay.read();
1922 let value = if let Some(v) =
1924 overlay.get(self.col as usize).and_then(|o| o.get(&self.key))
1925 {
1926 Ok(v.map(|i| i.as_ref().to_vec()))
1927 } else {
1928 let log = self.db.log.overlays();
1930 Ok(column.get(&self.key, log)?.map(|(v, _rc)| v))
1931 }?;
1932
1933 if let Some(data) = value {
1934 return unpack_node_data(data).map(|x| Some(x))
1935 }
1936
1937 return Ok(None)
1938 },
1939 Column::Tree(..) =>
1940 return Err(Error::InvalidConfiguration("Not a HashColumn.".to_string())),
1941 };
1942 }
1943
1944 fn get_node(&self, node_address: NodeAddress) -> Result<Option<(Vec<u8>, Children)>> {
1945 self.db.get_node(self.col, node_address, false)
1946 }
1947
1948 fn get_node_children(&self, node_address: NodeAddress) -> Result<Option<Children>> {
1949 self.db.get_node_children(self.col, node_address, false)
1950 }
1951}
1952
1953pub type IndexedCommitOverlay = HashMap<Key, (u64, Option<RcValue>), IdentityBuildHasher>;
1954pub type AddressCommitOverlay = HashMap<u64, (u64, RcValue)>;
1955pub type BTreeCommitOverlay = BTreeMap<RcKey, (u64, Option<RcValue>)>;
1956
1957#[derive(Debug)]
1958pub struct CommitOverlay {
1959 indexed: IndexedCommitOverlay,
1960 address: AddressCommitOverlay,
1961 btree_indexed: BTreeCommitOverlay,
1962}
1963
1964impl CommitOverlay {
1965 fn new() -> Self {
1966 CommitOverlay {
1967 indexed: Default::default(),
1968 address: Default::default(),
1969 btree_indexed: Default::default(),
1970 }
1971 }
1972
1973 #[cfg(test)]
1974 fn is_empty(&self) -> bool {
1975 self.indexed.is_empty() && self.address.is_empty() && self.btree_indexed.is_empty()
1976 }
1977}
1978
1979impl CommitOverlay {
1980 fn get_ref(&self, key: &[u8]) -> Option<Option<&RcValue>> {
1981 self.indexed.get(key).map(|(_, v)| v.as_ref())
1982 }
1983
1984 fn get(&self, key: &[u8]) -> Option<Option<RcValue>> {
1985 self.get_ref(key).map(|v| v.cloned())
1986 }
1987
1988 fn get_size(&self, key: &[u8]) -> Option<Option<u32>> {
1989 self.get_ref(key).map(|res| res.as_ref().map(|b| b.as_ref().len() as u32))
1990 }
1991
1992 fn get_address(&self, address: u64) -> Option<RcValue> {
1993 self.address.get(&address).map(|(_, v)| v.clone())
1994 }
1995
1996 fn btree_get(&self, key: &[u8]) -> Option<Option<&RcValue>> {
1997 self.btree_indexed.get(key).map(|(_, v)| v.as_ref())
1998 }
1999
2000 pub fn btree_next(&self, last_key: &crate::btree::LastKey) -> Option<(RcKey, Option<RcValue>)> {
2001 use crate::btree::LastKey;
2002 match &last_key {
2003 LastKey::Start => self
2004 .btree_indexed
2005 .range::<[u8], _>(..)
2006 .next()
2007 .map(|(k, (_, v))| (k.clone(), v.clone())),
2008 LastKey::End => None,
2009 LastKey::At(key) => self
2010 .btree_indexed
2011 .range::<[u8], _>((Bound::Excluded(key.as_slice()), Bound::Unbounded))
2012 .next()
2013 .map(|(k, (_, v))| (k.clone(), v.clone())),
2014 LastKey::Seeked(key) => self
2015 .btree_indexed
2016 .range::<[u8], _>((Bound::Included(key.as_slice()), Bound::Unbounded))
2017 .next()
2018 .map(|(k, (_, v))| (k.clone(), v.clone())),
2019 }
2020 }
2021
2022 pub fn btree_prev(&self, last_key: &crate::btree::LastKey) -> Option<(RcKey, Option<RcValue>)> {
2023 use crate::btree::LastKey;
2024 match &last_key {
2025 LastKey::End => self
2026 .btree_indexed
2027 .range::<[u8], _>(..)
2028 .rev()
2029 .next()
2030 .map(|(k, (_, v))| (k.clone(), v.clone())),
2031 LastKey::Start => None,
2032 LastKey::At(key) => self
2033 .btree_indexed
2034 .range::<[u8], _>((Bound::Unbounded, Bound::Excluded(key.as_slice())))
2035 .rev()
2036 .next()
2037 .map(|(k, (_, v))| (k.clone(), v.clone())),
2038 LastKey::Seeked(key) => self
2039 .btree_indexed
2040 .range::<[u8], _>((Bound::Unbounded, Bound::Included(key.as_slice())))
2041 .rev()
2042 .next()
2043 .map(|(k, (_, v))| (k.clone(), v.clone())),
2044 }
2045 }
2046}
2047
2048#[derive(Debug, PartialEq, Eq)]
2051pub enum Operation<Key, Value> {
2052 Set(Key, Value),
2054
2055 Dereference(Key),
2059
2060 Reference(Key),
2063
2064 InsertTree(Key, NewNode),
2066
2067 ReferenceTree(Key),
2069
2070 DereferenceTree(Key),
2073}
2074
2075impl<Key: Ord, Value: Eq> PartialOrd<Self> for Operation<Key, Value> {
2076 fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
2077 Some(self.cmp(other))
2078 }
2079}
2080
2081impl<Key: Ord, Value: Eq> Ord for Operation<Key, Value> {
2082 fn cmp(&self, other: &Self) -> std::cmp::Ordering {
2083 self.key().cmp(other.key())
2084 }
2085}
2086
2087impl<Key, Value> Operation<Key, Value> {
2088 pub fn key(&self) -> &Key {
2089 match self {
2090 Operation::Set(k, _) |
2091 Operation::Dereference(k) |
2092 Operation::Reference(k) |
2093 Operation::InsertTree(k, _) |
2094 Operation::ReferenceTree(k) |
2095 Operation::DereferenceTree(k) => k,
2096 }
2097 }
2098
2099 pub fn into_key(self) -> Key {
2100 match self {
2101 Operation::Set(k, _) |
2102 Operation::Dereference(k) |
2103 Operation::Reference(k) |
2104 Operation::InsertTree(k, _) |
2105 Operation::ReferenceTree(k) |
2106 Operation::DereferenceTree(k) => k,
2107 }
2108 }
2109}
2110
2111impl<K: AsRef<[u8]>, Value> Operation<K, Value> {
2112 pub fn to_key_vec(self) -> Operation<Vec<u8>, Value> {
2113 match self {
2114 Operation::Set(k, v) => Operation::Set(k.as_ref().to_vec(), v),
2115 Operation::Dereference(k) => Operation::Dereference(k.as_ref().to_vec()),
2116 Operation::Reference(k) => Operation::Reference(k.as_ref().to_vec()),
2117 Operation::InsertTree(k, n) => Operation::InsertTree(k.as_ref().to_vec(), n),
2118 Operation::ReferenceTree(k) => Operation::ReferenceTree(k.as_ref().to_vec()),
2119 Operation::DereferenceTree(k) => Operation::DereferenceTree(k.as_ref().to_vec()),
2120 }
2121 }
2122}
2123
2124#[derive(Debug, PartialEq, Eq)]
2125pub enum NodeChange {
2126 NewValue(u64, RcValue),
2128 IncrementReference(u64),
2130 DereferenceChildren(Vec<u8>, Key, Children),
2132}
2133
2134#[derive(Debug, Default)]
2135pub struct CommitChangeSet {
2136 pub indexed: HashMap<ColId, IndexedChangeSet>,
2137 pub btree_indexed: HashMap<ColId, BTreeChangeSet>,
2138 pub check_for_deferral: bool,
2139}
2140
2141#[derive(Debug)]
2142pub struct IndexedChangeSet {
2143 pub col: ColId,
2144 pub changes: Vec<Operation<Key, RcValue>>,
2145 pub node_changes: Vec<NodeChange>,
2146 pub used_trees: HashSet<Key>,
2147}
2148
2149impl IndexedChangeSet {
2150 pub fn new(col: ColId) -> Self {
2151 IndexedChangeSet {
2152 col,
2153 changes: Default::default(),
2154 node_changes: Default::default(),
2155 used_trees: Default::default(),
2156 }
2157 }
2158
2159 fn push<K: AsRef<[u8]>, V: Into<RcValue>>(
2160 &mut self,
2161 change: Operation<K, V>,
2162 options: &Options,
2163 db_version: u32,
2164 ) -> Result<()> {
2165 let salt = options.salt.unwrap_or_default();
2166 let hash_key = |key: &[u8]| -> Key {
2167 hash_key(key, &salt, options.columns[self.col as usize].uniform, db_version)
2168 };
2169
2170 self.push_change_hashed(match change {
2171 Operation::Set(k, v) => Operation::Set(hash_key(k.as_ref()), v.into()),
2172 Operation::Dereference(k) => Operation::Dereference(hash_key(k.as_ref())),
2173 Operation::Reference(k) => Operation::Reference(hash_key(k.as_ref())),
2174 Operation::InsertTree(..) |
2175 Operation::ReferenceTree(..) |
2176 Operation::DereferenceTree(..) =>
2177 return Err(Error::InvalidInput(format!(
2178 "Invalid operation for column {}",
2179 self.col
2180 ))),
2181 });
2182
2183 Ok(())
2184 }
2185
2186 fn push_change_hashed(&mut self, change: Operation<Key, RcValue>) {
2187 self.changes.push(change);
2188 }
2189
2190 fn push_node_change(&mut self, change: NodeChange) {
2191 self.node_changes.push(change);
2192 }
2193
2194 fn copy_to_overlay(
2195 &self,
2196 overlay: &mut CommitOverlay,
2197 record_id: u64,
2198 bytes: &mut usize,
2199 options: &Options,
2200 ) -> Result<()> {
2201 let ref_counted = options.columns[self.col as usize].ref_counted;
2202 for change in self.changes.iter() {
2203 match &change {
2204 Operation::Set(k, v) => {
2205 *bytes += k.len();
2206 *bytes += v.as_ref().len();
2207 overlay.indexed.insert(*k, (record_id, Some(v.clone())));
2208 },
2209 Operation::Dereference(k) => {
2210 if !ref_counted {
2212 overlay.indexed.insert(*k, (record_id, None));
2213 }
2214 },
2215 Operation::Reference(..) => {
2216 if !ref_counted {
2219 return Err(Error::InvalidInput(format!("No Rc for column {}", self.col)))
2220 }
2221 },
2222 Operation::InsertTree(..) |
2223 Operation::ReferenceTree(..) |
2224 Operation::DereferenceTree(..) =>
2225 return Err(Error::InvalidInput(format!(
2226 "Invalid operation for column {}",
2227 self.col
2228 ))),
2229 }
2230 }
2231 for change in self.node_changes.iter() {
2232 if let NodeChange::NewValue(address, val) = change {
2233 *bytes += val.as_ref().len();
2234 overlay.address.insert(*address, (record_id, val.clone()));
2235 }
2236 }
2237 Ok(())
2238 }
2239
2240 fn write_plan(
2241 &self,
2242 db: &Arc<DbInner>,
2243 col: ColId,
2244 column: &Column,
2245 writer: &mut crate::log::LogWriter,
2246 ops: &mut u64,
2247 reindex: &mut bool,
2248 ) -> Result<()> {
2249 let column = match column {
2250 Column::Hash(column) => column,
2251 Column::Tree(_) => {
2252 log::warn!(target: "parity-db", "Skipping unindex commit in indexed column");
2253 return Ok(())
2254 },
2255 };
2256 for change in self.changes.iter() {
2257 if let PlanOutcome::NeedReindex = column.write_plan(change, writer)? {
2258 *reindex = true;
2260 }
2261 *ops += 1;
2262 }
2263 for change in self.node_changes.iter() {
2264 match change {
2265 NodeChange::NewValue(address, val) => {
2266 column.write_address_value_plan(
2267 *address,
2268 val.clone(),
2269 false,
2270 val.as_ref().len() as u32,
2271 writer,
2272 )?;
2273 },
2274 NodeChange::IncrementReference(address) => {
2275 if let PlanOutcome::NeedReindex =
2276 column.write_address_inc_ref_plan(*address, writer)?
2277 {
2278 *reindex = true;
2279 }
2280 },
2281 NodeChange::DereferenceChildren(key, hash, children) => {
2282 if let Some((_root, rc)) = column.get(hash, writer)? {
2283 column.write_plan(&Operation::Dereference(*hash), writer)?;
2284 log::debug!(target: "parity-db", "Dereferencing root, rc={}", rc);
2285 if rc == 1 {
2286 let tree = db.get_tree(db, col, key, false).unwrap();
2287 if let Some(tree) = tree {
2288 let guard = tree.write();
2289 let mut num_removed = 0;
2290 self.write_dereference_children_plan(
2291 column,
2292 &guard,
2293 children,
2294 &mut num_removed,
2295 writer,
2296 )?;
2297 log::debug!(target: "parity-db", "Dereferenced tree {:?}, removed {}", &key[0..3], num_removed);
2298 }
2299 }
2300 }
2301 },
2303 }
2304 }
2305 Ok(())
2306 }
2307
2308 fn write_dereference_children_plan(
2309 &self,
2310 column: &HashColumn,
2311 guard: &RwLockWriteGuard<'_, Box<dyn TreeReader + Send + Sync>>,
2312 children: &Vec<u64>,
2313 num_removed: &mut u64,
2314 writer: &mut crate::log::LogWriter,
2315 ) -> Result<()> {
2316 for address in children {
2317 let node = guard.get_node_children(*address)?;
2321 let (remains, _outcome) = column.write_address_dec_ref_plan(*address, writer)?;
2322 if !remains {
2323 *num_removed += 1;
2325 if let Some(children) = node {
2326 self.write_dereference_children_plan(
2327 column,
2328 guard,
2329 &children,
2330 num_removed,
2331 writer,
2332 )?;
2333 } else {
2334 return Err(Error::InvalidConfiguration("Missing node data".to_string()))
2335 }
2336 }
2337 }
2338 Ok(())
2339 }
2340
2341 fn clean_overlay(&self, overlay: &mut CommitOverlay, record_id: u64) {
2342 use std::collections::hash_map::Entry;
2343 for change in self.changes.iter() {
2344 match change {
2345 Operation::Set(k, _) | Operation::Dereference(k) => {
2346 if let Entry::Occupied(e) = overlay.indexed.entry(*k) {
2347 if e.get().0 == record_id {
2348 e.remove_entry();
2349 }
2350 }
2351 },
2352 Operation::Reference(..) |
2353 Operation::InsertTree(..) |
2354 Operation::ReferenceTree(..) |
2355 Operation::DereferenceTree(..) => (),
2356 }
2357 }
2358 for change in self.node_changes.iter() {
2359 if let NodeChange::NewValue(address, _val) = change {
2360 if let Entry::Occupied(e) = overlay.address.entry(*address) {
2361 if e.get().0 == record_id {
2362 e.remove_entry();
2363 }
2364 }
2365 }
2366 }
2367 }
2368}
2369
2370pub mod check {
2372 pub enum CheckDisplay {
2374 None,
2376 Full,
2378 Short(u64),
2380 }
2381
2382 pub struct CheckOptions {
2384 pub column: Option<u8>,
2386 pub from: Option<u64>,
2388 pub bound: Option<u64>,
2390 pub display: CheckDisplay,
2392 pub fast: bool,
2394 pub validate_free_refs: bool,
2396 }
2397
2398 impl CheckOptions {
2399 pub fn new(
2401 column: Option<u8>,
2402 from: Option<u64>,
2403 bound: Option<u64>,
2404 display_content: bool,
2405 truncate_value_display: Option<u64>,
2406 fast: bool,
2407 validate_free_refs: bool,
2408 ) -> Self {
2409 let display = if display_content {
2410 match truncate_value_display {
2411 Some(t) => CheckDisplay::Short(t),
2412 None => CheckDisplay::Full,
2413 }
2414 } else {
2415 CheckDisplay::None
2416 };
2417 CheckOptions { column, from, bound, display, fast, validate_free_refs }
2418 }
2419 }
2420}
2421
2422#[derive(Eq, PartialEq, Clone, Copy)]
2423enum OpeningMode {
2424 Create,
2425 Write,
2426 ReadOnly,
2427}
2428
2429#[cfg(test)]
2430mod tests {
2431 use super::{Db, Options};
2432 use crate::{
2433 column::ColId,
2434 db::{DbInner, OpeningMode},
2435 ColumnOptions, Value,
2436 };
2437 use rand::Rng;
2438 use std::{
2439 collections::{BTreeMap, HashMap, HashSet},
2440 path::Path,
2441 };
2442 use tempfile::tempdir;
2443
2444 #[derive(Eq, PartialEq, Debug, Clone, Copy)]
2446 enum EnableCommitPipelineStages {
2447 #[allow(dead_code)]
2449 CommitOverlay,
2450 #[allow(dead_code)]
2452 LogOverlay,
2453 #[allow(dead_code)]
2455 DbFile,
2456 Standard,
2458 }
2459
2460 impl EnableCommitPipelineStages {
2461 fn options(&self, path: &Path, num_columns: u8) -> Options {
2462 Options {
2463 path: path.into(),
2464 sync_wal: true,
2465 sync_data: true,
2466 stats: true,
2467 salt: None,
2468 columns: (0..num_columns).map(|_| Default::default()).collect(),
2469 compression_threshold: HashMap::new(),
2470 with_background_thread: *self == Self::Standard,
2471 always_flush: *self == Self::DbFile,
2472 }
2473 }
2474
2475 fn run_stages(&self, db: &Db) {
2476 let db = &db.inner;
2477 if *self == EnableCommitPipelineStages::DbFile ||
2478 *self == EnableCommitPipelineStages::LogOverlay
2479 {
2480 while db.process_commits(db).unwrap() {}
2481 while db.process_reindex().unwrap() {}
2482 }
2483 if *self == EnableCommitPipelineStages::DbFile {
2484 let _ = db.log.flush_one(0).unwrap();
2485 while db.enact_logs(false).unwrap() {}
2486 let _ = db.clean_logs().unwrap();
2487 }
2488 }
2489
2490 fn check_empty_overlay(&self, db: &DbInner, col: ColId) -> bool {
2491 match self {
2492 EnableCommitPipelineStages::DbFile | EnableCommitPipelineStages::LogOverlay => {
2493 if let Some(overlay) = db.commit_overlay.read().get(col as usize) {
2494 if !overlay.is_empty() {
2495 let mut replayed = 5;
2496 while !overlay.is_empty() {
2497 if replayed > 0 {
2498 replayed -= 1;
2499 std::thread::sleep(std::time::Duration::from_millis(100));
2504 } else {
2505 return false
2506 }
2507 }
2508 }
2509 }
2510 },
2511 _ => (),
2512 }
2513 true
2514 }
2515 }
2516
2517 #[test]
2518 fn test_db_open_should_fail() {
2519 let tmp = tempdir().unwrap();
2520 let options = Options::with_columns(tmp.path(), 5);
2521 assert!(matches!(Db::open(&options), Err(crate::Error::DatabaseNotFound)));
2522 }
2523
2524 #[test]
2525 fn test_db_open_fail_then_recursively_create() {
2526 let tmp = tempdir().unwrap();
2527 let (db_path_first, db_path_last) = {
2528 let mut db_path_first = tmp.path().to_owned();
2529 db_path_first.push("nope");
2530
2531 let mut db_path_last = db_path_first.to_owned();
2532
2533 for p in ["does", "not", "yet", "exist"] {
2534 db_path_last.push(p);
2535 }
2536
2537 (db_path_first, db_path_last)
2538 };
2539
2540 assert!(
2541 !db_path_first.exists(),
2542 "That directory should not have existed at this point (dir: {db_path_first:?})"
2543 );
2544
2545 let options = Options::with_columns(&db_path_last, 5);
2546 assert!(matches!(Db::open(&options), Err(crate::Error::DatabaseNotFound)));
2547
2548 assert!(!db_path_first.exists(), "That directory should remain non-existent. Did the `open(create: false)` nonetheless create a directory? (dir: {db_path_first:?})");
2549 assert!(Db::open_or_create(&options).is_ok(), "New database should be created");
2550
2551 assert!(
2552 db_path_first.is_dir(),
2553 "A directory should have been been created (dir: {db_path_first:?})"
2554 );
2555 assert!(
2556 db_path_last.is_dir(),
2557 "A directory should have been been created (dir: {db_path_last:?})"
2558 );
2559 }
2560
2561 #[test]
2562 fn test_db_open_or_create() {
2563 let tmp = tempdir().unwrap();
2564 let options = Options::with_columns(tmp.path(), 5);
2565 assert!(Db::open_or_create(&options).is_ok(), "New database should be created");
2566 assert!(Db::open(&options).is_ok(), "Existing database should be reopened");
2567 }
2568
2569 #[test]
2570 fn test_indexed_keyvalues() {
2571 test_indexed_keyvalues_inner(EnableCommitPipelineStages::CommitOverlay);
2572 test_indexed_keyvalues_inner(EnableCommitPipelineStages::LogOverlay);
2573 test_indexed_keyvalues_inner(EnableCommitPipelineStages::DbFile);
2574 test_indexed_keyvalues_inner(EnableCommitPipelineStages::Standard);
2575 }
2576 fn test_indexed_keyvalues_inner(db_test: EnableCommitPipelineStages) {
2577 let tmp = tempdir().unwrap();
2578 let options = db_test.options(tmp.path(), 5);
2579 let col_nb = 0;
2580
2581 let key1 = b"key1".to_vec();
2582 let key2 = b"key2".to_vec();
2583 let key3 = b"key3".to_vec();
2584
2585 let db = Db::open_inner(&options, OpeningMode::Create).unwrap();
2586 assert!(db.get(col_nb, key1.as_slice()).unwrap().is_none());
2587
2588 db.commit(vec![(col_nb, key1.clone(), Some(b"value1".to_vec()))]).unwrap();
2589 db_test.run_stages(&db);
2590 assert!(db_test.check_empty_overlay(&db.inner, col_nb));
2591
2592 assert_eq!(db.get(col_nb, key1.as_slice()).unwrap(), Some(b"value1".to_vec()));
2593
2594 db.commit(vec![
2595 (col_nb, key1.clone(), None),
2596 (col_nb, key2.clone(), Some(b"value2".to_vec())),
2597 (col_nb, key3.clone(), Some(b"value3".to_vec())),
2598 ])
2599 .unwrap();
2600 db_test.run_stages(&db);
2601 assert!(db_test.check_empty_overlay(&db.inner, col_nb));
2602
2603 assert!(db.get(col_nb, key1.as_slice()).unwrap().is_none());
2604 assert_eq!(db.get(col_nb, key2.as_slice()).unwrap(), Some(b"value2".to_vec()));
2605 assert_eq!(db.get(col_nb, key3.as_slice()).unwrap(), Some(b"value3".to_vec()));
2606
2607 db.commit(vec![
2608 (col_nb, key2.clone(), Some(b"value2b".to_vec())),
2609 (col_nb, key3.clone(), None),
2610 ])
2611 .unwrap();
2612 db_test.run_stages(&db);
2613 assert!(db_test.check_empty_overlay(&db.inner, col_nb));
2614
2615 assert!(db.get(col_nb, key1.as_slice()).unwrap().is_none());
2616 assert_eq!(db.get(col_nb, key2.as_slice()).unwrap(), Some(b"value2b".to_vec()));
2617 assert_eq!(db.get(col_nb, key3.as_slice()).unwrap(), None);
2618 }
2619
2620 #[test]
2621 fn test_indexed_overlay_against_backend() {
2622 let tmp = tempdir().unwrap();
2623 let db_test = EnableCommitPipelineStages::DbFile;
2624 let options = db_test.options(tmp.path(), 5);
2625 let col_nb = 0;
2626
2627 let key1 = b"key1".to_vec();
2628 let key2 = b"key2".to_vec();
2629 let key3 = b"key3".to_vec();
2630
2631 let db = Db::open_inner(&options, OpeningMode::Create).unwrap();
2632
2633 db.commit(vec![
2634 (col_nb, key1.clone(), Some(b"value1".to_vec())),
2635 (col_nb, key2.clone(), Some(b"value2".to_vec())),
2636 (col_nb, key3.clone(), Some(b"value3".to_vec())),
2637 ])
2638 .unwrap();
2639 db_test.run_stages(&db);
2640 drop(db);
2641
2642 std::thread::sleep(std::time::Duration::from_millis(100));
2644
2645 let db_test = EnableCommitPipelineStages::CommitOverlay;
2646 let options = db_test.options(tmp.path(), 5);
2647 let db = Db::open_inner(&options, OpeningMode::Write).unwrap();
2648 assert_eq!(db.get(col_nb, key1.as_slice()).unwrap(), Some(b"value1".to_vec()));
2649 assert_eq!(db.get(col_nb, key2.as_slice()).unwrap(), Some(b"value2".to_vec()));
2650 assert_eq!(db.get(col_nb, key3.as_slice()).unwrap(), Some(b"value3".to_vec()));
2651 db.commit(vec![
2652 (col_nb, key2.clone(), Some(b"value2b".to_vec())),
2653 (col_nb, key3.clone(), None),
2654 ])
2655 .unwrap();
2656 db_test.run_stages(&db);
2657
2658 assert_eq!(db.get(col_nb, key1.as_slice()).unwrap(), Some(b"value1".to_vec()));
2659 assert_eq!(db.get(col_nb, key2.as_slice()).unwrap(), Some(b"value2b".to_vec()));
2660 assert_eq!(db.get(col_nb, key3.as_slice()).unwrap(), None);
2661 }
2662
2663 #[test]
2664 fn test_add_column() {
2665 let tmp = tempdir().unwrap();
2666 let db_test = EnableCommitPipelineStages::DbFile;
2667 let mut options = db_test.options(tmp.path(), 1);
2668 options.salt = Some(options.salt.unwrap_or_default());
2669
2670 let old_col_id = 0;
2671 let new_col_id = 1;
2672 let new_col_indexed_id = 2;
2673
2674 let key1 = b"key1".to_vec();
2675 let key2 = b"key2".to_vec();
2676 let key3 = b"key3".to_vec();
2677
2678 let db = Db::open_inner(&options, OpeningMode::Create).unwrap();
2679
2680 db.commit(vec![
2681 (old_col_id, key1.clone(), Some(b"value1".to_vec())),
2682 (old_col_id, key2.clone(), Some(b"value2".to_vec())),
2683 (old_col_id, key3.clone(), Some(b"value3".to_vec())),
2684 ])
2685 .unwrap();
2686 db_test.run_stages(&db);
2687
2688 drop(db);
2689
2690 Db::add_column(&mut options, ColumnOptions { btree_index: false, ..Default::default() })
2691 .unwrap();
2692
2693 Db::add_column(&mut options, ColumnOptions { btree_index: true, ..Default::default() })
2694 .unwrap();
2695
2696 let mut options = db_test.options(tmp.path(), 3);
2697 options.columns[new_col_indexed_id as usize].btree_index = true;
2698
2699 let db_test = EnableCommitPipelineStages::DbFile;
2700 let db = Db::open_inner(&options, OpeningMode::Create).unwrap();
2701
2702 assert_eq!(db.num_columns(), 3);
2704
2705 let new_key1 = b"abcdef".to_vec();
2706 let new_key2 = b"123456".to_vec();
2707
2708 db.commit(vec![
2710 (new_col_id, new_key1.clone(), Some(new_key1.to_vec())),
2711 (new_col_id, new_key2.clone(), Some(new_key2.to_vec())),
2712 (new_col_indexed_id, new_key1.clone(), Some(new_key1.to_vec())),
2713 (new_col_indexed_id, new_key2.clone(), Some(new_key2.to_vec())),
2714 ])
2715 .unwrap();
2716 db_test.run_stages(&db);
2717
2718 drop(db);
2719
2720 let db = Db::open_inner(&options, OpeningMode::Create).unwrap();
2722
2723 assert_eq!(db.get(old_col_id, key1.as_slice()).unwrap(), Some(b"value1".to_vec()));
2724 assert_eq!(db.get(old_col_id, key2.as_slice()).unwrap(), Some(b"value2".to_vec()));
2725 assert_eq!(db.get(old_col_id, key3.as_slice()).unwrap(), Some(b"value3".to_vec()));
2726
2727 assert_eq!(db.get(new_col_id, new_key1.as_slice()).unwrap(), Some(new_key1.to_vec()));
2729 assert_eq!(db.get(new_col_id, new_key2.as_slice()).unwrap(), Some(new_key2.to_vec()));
2730 assert_eq!(
2731 db.get(new_col_indexed_id, new_key1.as_slice()).unwrap(),
2732 Some(new_key1.to_vec())
2733 );
2734 assert_eq!(
2735 db.get(new_col_indexed_id, new_key2.as_slice()).unwrap(),
2736 Some(new_key2.to_vec())
2737 );
2738 }
2739
2740 #[test]
2741 fn test_indexed_btree_1() {
2742 test_indexed_btree_inner(EnableCommitPipelineStages::CommitOverlay, false);
2743 test_indexed_btree_inner(EnableCommitPipelineStages::LogOverlay, false);
2744 test_indexed_btree_inner(EnableCommitPipelineStages::DbFile, false);
2745 test_indexed_btree_inner(EnableCommitPipelineStages::Standard, false);
2746 test_indexed_btree_inner(EnableCommitPipelineStages::CommitOverlay, true);
2747 test_indexed_btree_inner(EnableCommitPipelineStages::LogOverlay, true);
2748 test_indexed_btree_inner(EnableCommitPipelineStages::DbFile, true);
2749 test_indexed_btree_inner(EnableCommitPipelineStages::Standard, true);
2750 }
2751 fn test_indexed_btree_inner(db_test: EnableCommitPipelineStages, long_key: bool) {
2752 let tmp = tempdir().unwrap();
2753 let col_nb = 0u8;
2754 let mut options = db_test.options(tmp.path(), 5);
2755 options.columns[col_nb as usize].btree_index = true;
2756
2757 let (key1, key2, key3, key4) = if !long_key {
2758 (b"key1".to_vec(), b"key2".to_vec(), b"key3".to_vec(), b"key4".to_vec())
2759 } else {
2760 let key2 = vec![2; 272];
2761 let mut key3 = key2.clone();
2762 key3[271] = 3;
2763 (vec![1; 953], key2, key3, vec![4; 79])
2764 };
2765
2766 let db = Db::open_inner(&options, OpeningMode::Create).unwrap();
2767 assert_eq!(db.get(col_nb, &key1).unwrap(), None);
2768
2769 let mut iter = db.iter(col_nb).unwrap();
2770 assert_eq!(iter.next().unwrap(), None);
2771 assert_eq!(iter.prev().unwrap(), None);
2772
2773 db.commit(vec![(col_nb, key1.clone(), Some(b"value1".to_vec()))]).unwrap();
2774 db_test.run_stages(&db);
2775
2776 assert_eq!(db.get(col_nb, &key1).unwrap(), Some(b"value1".to_vec()));
2777 iter.seek_to_first().unwrap();
2778 assert_eq!(iter.next().unwrap(), Some((key1.clone(), b"value1".to_vec())));
2779 assert_eq!(iter.next().unwrap(), None);
2780 assert_eq!(iter.prev().unwrap(), Some((key1.clone(), b"value1".to_vec())));
2781 assert_eq!(iter.prev().unwrap(), None);
2782 assert_eq!(iter.next().unwrap(), Some((key1.clone(), b"value1".to_vec())));
2783 assert_eq!(iter.next().unwrap(), None);
2784
2785 iter.seek_to_first().unwrap();
2786 assert_eq!(iter.next().unwrap(), Some((key1.clone(), b"value1".to_vec())));
2787 assert_eq!(iter.prev().unwrap(), None);
2788
2789 iter.seek(&[0xff]).unwrap();
2790 assert_eq!(iter.prev().unwrap(), Some((key1.clone(), b"value1".to_vec())));
2791 assert_eq!(iter.prev().unwrap(), None);
2792
2793 db.commit(vec![
2794 (col_nb, key1.clone(), None),
2795 (col_nb, key2.clone(), Some(b"value2".to_vec())),
2796 (col_nb, key3.clone(), Some(b"value3".to_vec())),
2797 ])
2798 .unwrap();
2799 db_test.run_stages(&db);
2800
2801 assert_eq!(db.get(col_nb, &key1).unwrap(), None);
2802 assert_eq!(db.get(col_nb, &key2).unwrap(), Some(b"value2".to_vec()));
2803 assert_eq!(db.get(col_nb, &key3).unwrap(), Some(b"value3".to_vec()));
2804
2805 iter.seek(key2.as_slice()).unwrap();
2806 assert_eq!(iter.next().unwrap(), Some((key2.clone(), b"value2".to_vec())));
2807 assert_eq!(iter.next().unwrap(), Some((key3.clone(), b"value3".to_vec())));
2808 assert_eq!(iter.next().unwrap(), None);
2809
2810 iter.seek(key3.as_slice()).unwrap();
2811 assert_eq!(iter.prev().unwrap(), Some((key3.clone(), b"value3".to_vec())));
2812 assert_eq!(iter.prev().unwrap(), Some((key2.clone(), b"value2".to_vec())));
2813 assert_eq!(iter.prev().unwrap(), None);
2814
2815 db.commit(vec![
2816 (col_nb, key2.clone(), Some(b"value2b".to_vec())),
2817 (col_nb, key4.clone(), Some(b"value4".to_vec())),
2818 (col_nb, key3.clone(), None),
2819 ])
2820 .unwrap();
2821 db_test.run_stages(&db);
2822
2823 assert_eq!(db.get(col_nb, &key1).unwrap(), None);
2824 assert_eq!(db.get(col_nb, &key3).unwrap(), None);
2825 assert_eq!(db.get(col_nb, &key2).unwrap(), Some(b"value2b".to_vec()));
2826 assert_eq!(db.get(col_nb, &key4).unwrap(), Some(b"value4".to_vec()));
2827 let mut key22 = key2.clone();
2828 key22.push(2);
2829 iter.seek(key22.as_slice()).unwrap();
2830 assert_eq!(iter.next().unwrap(), Some((key4, b"value4".to_vec())));
2831 assert_eq!(iter.next().unwrap(), None);
2832 }
2833
2834 #[test]
2835 fn test_indexed_btree_2() {
2836 test_indexed_btree_inner_2(EnableCommitPipelineStages::CommitOverlay);
2837 test_indexed_btree_inner_2(EnableCommitPipelineStages::LogOverlay);
2838 }
2839 fn test_indexed_btree_inner_2(db_test: EnableCommitPipelineStages) {
2840 let tmp = tempdir().unwrap();
2841 let col_nb = 0u8;
2842 let mut options = db_test.options(tmp.path(), 5);
2843 options.columns[col_nb as usize].btree_index = true;
2844
2845 let key1 = b"key1".to_vec();
2846 let key2 = b"key2".to_vec();
2847 let key3 = b"key3".to_vec();
2848
2849 let db = Db::open_inner(&options, OpeningMode::Create).unwrap();
2850 let mut iter = db.iter(col_nb).unwrap();
2851 assert_eq!(db.get(col_nb, &key1).unwrap(), None);
2852 assert_eq!(iter.next().unwrap(), None);
2853
2854 db.commit(vec![(col_nb, key1.clone(), Some(b"value1".to_vec()))]).unwrap();
2855 EnableCommitPipelineStages::DbFile.run_stages(&db);
2856 drop(db);
2857
2858 std::thread::sleep(std::time::Duration::from_millis(100));
2860
2861 let db = Db::open_inner(&options, OpeningMode::Write).unwrap();
2862
2863 let mut iter = db.iter(col_nb).unwrap();
2864 assert_eq!(db.get(col_nb, &key1).unwrap(), Some(b"value1".to_vec()));
2865 iter.seek_to_first().unwrap();
2866 assert_eq!(iter.next().unwrap(), Some((key1.clone(), b"value1".to_vec())));
2867 assert_eq!(iter.next().unwrap(), None);
2868
2869 db.commit(vec![
2870 (col_nb, key1.clone(), None),
2871 (col_nb, key2.clone(), Some(b"value2".to_vec())),
2872 (col_nb, key3.clone(), Some(b"value3".to_vec())),
2873 ])
2874 .unwrap();
2875 db_test.run_stages(&db);
2876
2877 assert_eq!(db.get(col_nb, &key1).unwrap(), None);
2878 assert_eq!(db.get(col_nb, &key2).unwrap(), Some(b"value2".to_vec()));
2879 assert_eq!(db.get(col_nb, &key3).unwrap(), Some(b"value3".to_vec()));
2880 iter.seek(key2.as_slice()).unwrap();
2881 assert_eq!(iter.next().unwrap(), Some((key2.clone(), b"value2".to_vec())));
2882 assert_eq!(iter.next().unwrap(), Some((key3.clone(), b"value3".to_vec())));
2883 assert_eq!(iter.next().unwrap(), None);
2884
2885 iter.seek_to_last().unwrap();
2886 assert_eq!(iter.prev().unwrap(), Some((key3, b"value3".to_vec())));
2887 assert_eq!(iter.prev().unwrap(), Some((key2.clone(), b"value2".to_vec())));
2888 assert_eq!(iter.prev().unwrap(), None);
2889 }
2890
2891 #[test]
2892 fn test_indexed_btree_3() {
2893 test_indexed_btree_inner_3(EnableCommitPipelineStages::CommitOverlay);
2894 test_indexed_btree_inner_3(EnableCommitPipelineStages::LogOverlay);
2895 test_indexed_btree_inner_3(EnableCommitPipelineStages::DbFile);
2896 test_indexed_btree_inner_3(EnableCommitPipelineStages::Standard);
2897 }
2898
2899 fn test_indexed_btree_inner_3(db_test: EnableCommitPipelineStages) {
2900 use rand::SeedableRng;
2901
2902 use std::collections::BTreeSet;
2903
2904 let mut rng = rand::rngs::SmallRng::seed_from_u64(0);
2905
2906 let tmp = tempdir().unwrap();
2907 let col_nb = 0u8;
2908 let mut options = db_test.options(tmp.path(), 5);
2909 options.columns[col_nb as usize].btree_index = true;
2910
2911 let db = Db::open_inner(&options, OpeningMode::Create).unwrap();
2912
2913 db.commit(
2914 (0u64..1024)
2915 .map(|i| (0, i.to_be_bytes().to_vec(), Some(i.to_be_bytes().to_vec())))
2916 .chain((0u64..1024).step_by(2).map(|i| (0, i.to_be_bytes().to_vec(), None))),
2917 )
2918 .unwrap();
2919 let expected = (0u64..1024).filter(|i| i % 2 == 1).collect::<BTreeSet<_>>();
2920 let mut iter = db.iter(0).unwrap();
2921
2922 for _ in 0..100 {
2923 let at = rng.random_range(0u64..=1024);
2924 iter.seek(&at.to_be_bytes()).unwrap();
2925
2926 let mut prev_run: bool = rng.random();
2927 let at = if prev_run {
2928 let take = rng.random_range(1..100);
2929 let got = std::iter::from_fn(|| iter.next().unwrap())
2930 .map(|(k, _)| u64::from_be_bytes(k.try_into().unwrap()))
2931 .take(take)
2932 .collect::<Vec<_>>();
2933 let expected = expected.range(at..).take(take).copied().collect::<Vec<_>>();
2934 assert_eq!(got, expected);
2935 if got.is_empty() {
2936 prev_run = false;
2937 }
2938 if got.len() < take {
2939 prev_run = false;
2940 }
2941 expected.last().copied().unwrap_or(at)
2942 } else {
2943 at
2944 };
2945
2946 let at = {
2947 let take = rng.random_range(1..100);
2948 let got = std::iter::from_fn(|| iter.prev().unwrap())
2949 .map(|(k, _)| u64::from_be_bytes(k.try_into().unwrap()))
2950 .take(take)
2951 .collect::<Vec<_>>();
2952 let expected = if prev_run {
2953 expected.range(..at).rev().take(take).copied().collect::<Vec<_>>()
2954 } else {
2955 expected.range(..=at).rev().take(take).copied().collect::<Vec<_>>()
2956 };
2957 assert_eq!(got, expected);
2958 prev_run = !got.is_empty();
2959 if take > got.len() {
2960 prev_run = false;
2961 }
2962 expected.last().copied().unwrap_or(at)
2963 };
2964
2965 let take = rng.random_range(1..100);
2966 let mut got = std::iter::from_fn(|| iter.next().unwrap())
2967 .map(|(k, _)| u64::from_be_bytes(k.try_into().unwrap()))
2968 .take(take)
2969 .collect::<Vec<_>>();
2970 let mut expected = expected.range(at..).take(take).copied().collect::<Vec<_>>();
2971 if prev_run {
2972 expected = expected.split_off(1);
2973 if got.len() == take {
2974 got.pop();
2975 }
2976 }
2977 assert_eq!(got, expected);
2978 }
2979
2980 let take = rng.random_range(20..100);
2981 iter.seek_to_last().unwrap();
2982 let got = std::iter::from_fn(|| iter.prev().unwrap())
2983 .map(|(k, _)| u64::from_be_bytes(k.try_into().unwrap()))
2984 .take(take)
2985 .collect::<Vec<_>>();
2986 let expected = expected.iter().rev().take(take).copied().collect::<Vec<_>>();
2987 assert_eq!(got, expected);
2988 }
2989
2990 fn test_basic(change_set: &[(Vec<u8>, Option<Vec<u8>>)]) {
2991 test_basic_inner(change_set, false, false);
2992 test_basic_inner(change_set, false, true);
2993 test_basic_inner(change_set, true, false);
2994 test_basic_inner(change_set, true, true);
2995 }
2996
2997 fn test_basic_inner(
2998 change_set: &[(Vec<u8>, Option<Vec<u8>>)],
2999 btree_index: bool,
3000 ref_counted: bool,
3001 ) {
3002 let tmp = tempdir().unwrap();
3003 let col_nb = 1u8;
3004 let db_test = EnableCommitPipelineStages::DbFile;
3005 let mut options = db_test.options(tmp.path(), 2);
3006 options.columns[col_nb as usize].btree_index = btree_index;
3007 options.columns[col_nb as usize].ref_counted = ref_counted;
3008 options.columns[col_nb as usize].preimage = ref_counted;
3009 assert!(!(ref_counted && db_test == EnableCommitPipelineStages::CommitOverlay));
3011 let db = Db::open_inner(&options, OpeningMode::Create).unwrap();
3012
3013 let iter = btree_index.then(|| db.iter(col_nb).unwrap());
3014 assert_eq!(iter.and_then(|mut i| i.next().unwrap()), None);
3015
3016 db.commit(change_set.iter().map(|(k, v)| (col_nb, k.clone(), v.clone())))
3017 .unwrap();
3018 db_test.run_stages(&db);
3019
3020 let mut keys = HashSet::new();
3021 let mut expected_count: u64 = 0;
3022 for (k, v) in change_set.iter() {
3023 if v.is_some() {
3024 if keys.insert(k) {
3025 expected_count += 1;
3026 }
3027 } else if keys.remove(k) {
3028 expected_count -= 1;
3029 }
3030 }
3031 if ref_counted {
3032 let mut state: BTreeMap<Vec<u8>, Option<(Vec<u8>, usize)>> = Default::default();
3033 for (k, v) in change_set.iter() {
3034 let mut remove = false;
3035 let mut insert = false;
3036 match state.get_mut(k) {
3037 Some(Some((_, counter))) =>
3038 if v.is_some() {
3039 *counter += 1;
3040 } else if *counter == 1 {
3041 remove = true;
3042 } else {
3043 *counter -= 1;
3044 },
3045 Some(None) | None =>
3046 if v.is_some() {
3047 insert = true;
3048 },
3049 }
3050 if insert {
3051 state.insert(k.clone(), v.clone().map(|v| (v, 1)));
3052 }
3053 if remove {
3054 state.remove(k);
3055 }
3056 }
3057 for (key, value) in state {
3058 assert_eq!(db.get(col_nb, &key).unwrap(), value.map(|v| v.0));
3059 }
3060 } else {
3061 let stats = db.stats();
3062 if let Some(stats) = stats.columns[col_nb as usize].as_ref() {
3064 assert_eq!(stats.total_values, expected_count);
3065 }
3066
3067 let state: BTreeMap<Vec<u8>, Option<Vec<u8>>> =
3068 change_set.iter().map(|(k, v)| (k.clone(), v.clone())).collect();
3069 for (key, value) in state.iter() {
3070 assert_eq!(&db.get(col_nb, key).unwrap(), value);
3071 }
3072 }
3073 }
3074
3075 #[test]
3076 fn test_random() {
3077 fdlimit::raise_fd_limit().unwrap();
3078 for i in 0..100 {
3079 test_random_inner(60, 60, i);
3080 }
3081 for i in 0..50 {
3082 test_random_inner(20, 60, i);
3083 }
3084 }
3085 fn test_random_inner(size: usize, key_size: usize, seed: u64) {
3086 use rand::{RngCore, SeedableRng};
3087 let mut rng = rand::rngs::SmallRng::seed_from_u64(seed);
3088 let mut data = Vec::<(Vec<u8>, Option<Vec<u8>>)>::new();
3089 for i in 0..size {
3090 let nb_delete: u32 = rng.next_u32(); let nb_delete = (nb_delete as usize % size) / 2;
3092 let mut key = vec![0u8; key_size];
3093 rng.fill_bytes(&mut key[..]);
3094 let value = if i > size - nb_delete {
3095 let random_key = rng.next_u32();
3096 let random_key = (random_key % 4) > 0;
3097 if !random_key {
3098 key = data[i - size / 2].0.clone();
3099 }
3100 None
3101 } else {
3102 Some(key.clone())
3103 };
3104 let var_keysize = rng.next_u32();
3105 let var_keysize = var_keysize as usize % (key_size / 2);
3106 key.truncate(key_size - var_keysize);
3107 data.push((key, value));
3108 }
3109 test_basic(&data[..]);
3110 }
3111
3112 #[test]
3113 fn test_simple() {
3114 test_basic(&[
3115 (b"key1".to_vec(), Some(b"value1".to_vec())),
3116 (b"key1".to_vec(), Some(b"value1".to_vec())),
3117 (b"key1".to_vec(), None),
3118 ]);
3119 test_basic(&[
3120 (b"key1".to_vec(), Some(b"value1".to_vec())),
3121 (b"key1".to_vec(), Some(b"value1".to_vec())),
3122 (b"key1".to_vec(), None),
3123 (b"key1".to_vec(), None),
3124 ]);
3125 test_basic(&[
3126 (b"key1".to_vec(), Some(b"value1".to_vec())),
3127 (b"key1".to_vec(), Some(b"value2".to_vec())),
3128 ]);
3129 test_basic(&[(b"key1".to_vec(), Some(b"value1".to_vec()))]);
3130 test_basic(&[
3131 (b"key1".to_vec(), Some(b"value1".to_vec())),
3132 (b"key2".to_vec(), Some(b"value2".to_vec())),
3133 ]);
3134 test_basic(&[
3135 (b"key1".to_vec(), Some(b"value1".to_vec())),
3136 (b"key2".to_vec(), Some(b"value2".to_vec())),
3137 (b"key3".to_vec(), Some(b"value3".to_vec())),
3138 ]);
3139 test_basic(&[
3140 (b"key1".to_vec(), Some(b"value1".to_vec())),
3141 (b"key3".to_vec(), Some(b"value3".to_vec())),
3142 (b"key2".to_vec(), Some(b"value2".to_vec())),
3143 ]);
3144 test_basic(&[
3145 (b"key3".to_vec(), Some(b"value3".to_vec())),
3146 (b"key2".to_vec(), Some(b"value2".to_vec())),
3147 (b"key1".to_vec(), Some(b"value1".to_vec())),
3148 ]);
3149 test_basic(&[
3150 (b"key1".to_vec(), Some(b"value1".to_vec())),
3151 (b"key2".to_vec(), Some(b"value2".to_vec())),
3152 (b"key3".to_vec(), Some(b"value3".to_vec())),
3153 (b"key4".to_vec(), Some(b"value4".to_vec())),
3154 ]);
3155 test_basic(&[
3156 (b"key1".to_vec(), Some(b"value1".to_vec())),
3157 (b"key2".to_vec(), Some(b"value2".to_vec())),
3158 (b"key3".to_vec(), Some(b"value3".to_vec())),
3159 (b"key4".to_vec(), Some(b"value4".to_vec())),
3160 (b"key5".to_vec(), Some(b"value5".to_vec())),
3161 ]);
3162 test_basic(&[
3163 (b"key5".to_vec(), Some(b"value5".to_vec())),
3164 (b"key3".to_vec(), Some(b"value3".to_vec())),
3165 (b"key4".to_vec(), Some(b"value4".to_vec())),
3166 (b"key2".to_vec(), Some(b"value2".to_vec())),
3167 (b"key1".to_vec(), Some(b"value1".to_vec())),
3168 ]);
3169 test_basic(&[
3170 (b"key5".to_vec(), Some(b"value5".to_vec())),
3171 (b"key3".to_vec(), Some(b"value3".to_vec())),
3172 (b"key4".to_vec(), Some(b"value4".to_vec())),
3173 (b"key2".to_vec(), Some(b"value2".to_vec())),
3174 (b"key1".to_vec(), Some(b"value1".to_vec())),
3175 (b"key11".to_vec(), Some(b"value31".to_vec())),
3176 (b"key12".to_vec(), Some(b"value32".to_vec())),
3177 ]);
3178 test_basic(&[
3179 (b"key5".to_vec(), Some(b"value5".to_vec())),
3180 (b"key3".to_vec(), Some(b"value3".to_vec())),
3181 (b"key4".to_vec(), Some(b"value4".to_vec())),
3182 (b"key2".to_vec(), Some(b"value2".to_vec())),
3183 (b"key1".to_vec(), Some(b"value1".to_vec())),
3184 (b"key51".to_vec(), Some(b"value31".to_vec())),
3185 (b"key52".to_vec(), Some(b"value32".to_vec())),
3186 ]);
3187 test_basic(&[
3188 (b"key5".to_vec(), Some(b"value5".to_vec())),
3189 (b"key3".to_vec(), Some(b"value3".to_vec())),
3190 (b"key4".to_vec(), Some(b"value4".to_vec())),
3191 (b"key2".to_vec(), Some(b"value2".to_vec())),
3192 (b"key1".to_vec(), Some(b"value1".to_vec())),
3193 (b"key31".to_vec(), Some(b"value31".to_vec())),
3194 (b"key32".to_vec(), Some(b"value32".to_vec())),
3195 ]);
3196 test_basic(&[
3197 (b"key1".to_vec(), Some(b"value5".to_vec())),
3198 (b"key2".to_vec(), Some(b"value3".to_vec())),
3199 (b"key3".to_vec(), Some(b"value4".to_vec())),
3200 (b"key4".to_vec(), Some(b"value7".to_vec())),
3201 (b"key5".to_vec(), Some(b"value2".to_vec())),
3202 (b"key6".to_vec(), Some(b"value1".to_vec())),
3203 (b"key3".to_vec(), None),
3204 ]);
3205 test_basic(&[
3206 (b"key1".to_vec(), Some(b"value5".to_vec())),
3207 (b"key2".to_vec(), Some(b"value3".to_vec())),
3208 (b"key3".to_vec(), Some(b"value4".to_vec())),
3209 (b"key4".to_vec(), Some(b"value7".to_vec())),
3210 (b"key5".to_vec(), Some(b"value2".to_vec())),
3211 (b"key0".to_vec(), Some(b"value1".to_vec())),
3212 (b"key3".to_vec(), None),
3213 ]);
3214 test_basic(&[
3215 (b"key1".to_vec(), Some(b"value5".to_vec())),
3216 (b"key2".to_vec(), Some(b"value3".to_vec())),
3217 (b"key3".to_vec(), Some(b"value4".to_vec())),
3218 (b"key4".to_vec(), Some(b"value7".to_vec())),
3219 (b"key5".to_vec(), Some(b"value2".to_vec())),
3220 (b"key3".to_vec(), None),
3221 ]);
3222 test_basic(&[
3223 (b"key1".to_vec(), Some(b"value5".to_vec())),
3224 (b"key4".to_vec(), Some(b"value3".to_vec())),
3225 (b"key5".to_vec(), Some(b"value4".to_vec())),
3226 (b"key6".to_vec(), Some(b"value4".to_vec())),
3227 (b"key7".to_vec(), Some(b"value2".to_vec())),
3228 (b"key8".to_vec(), Some(b"value1".to_vec())),
3229 (b"key5".to_vec(), None),
3230 ]);
3231 test_basic(&[
3232 (b"key1".to_vec(), Some(b"value5".to_vec())),
3233 (b"key4".to_vec(), Some(b"value3".to_vec())),
3234 (b"key5".to_vec(), Some(b"value4".to_vec())),
3235 (b"key7".to_vec(), Some(b"value2".to_vec())),
3236 (b"key8".to_vec(), Some(b"value1".to_vec())),
3237 (b"key3".to_vec(), None),
3238 ]);
3239 test_basic(&[
3240 (b"key5".to_vec(), Some(b"value5".to_vec())),
3241 (b"key3".to_vec(), Some(b"value3".to_vec())),
3242 (b"key4".to_vec(), Some(b"value4".to_vec())),
3243 (b"key2".to_vec(), Some(b"value2".to_vec())),
3244 (b"key1".to_vec(), Some(b"value1".to_vec())),
3245 (b"key5".to_vec(), None),
3246 (b"key3".to_vec(), None),
3247 ]);
3248 test_basic(&[
3249 (b"key5".to_vec(), Some(b"value5".to_vec())),
3250 (b"key3".to_vec(), Some(b"value3".to_vec())),
3251 (b"key4".to_vec(), Some(b"value4".to_vec())),
3252 (b"key2".to_vec(), Some(b"value2".to_vec())),
3253 (b"key1".to_vec(), Some(b"value1".to_vec())),
3254 (b"key5".to_vec(), None),
3255 (b"key3".to_vec(), None),
3256 (b"key2".to_vec(), None),
3257 (b"key4".to_vec(), None),
3258 ]);
3259 test_basic(&[
3260 (b"key5".to_vec(), Some(b"value5".to_vec())),
3261 (b"key3".to_vec(), Some(b"value3".to_vec())),
3262 (b"key4".to_vec(), Some(b"value4".to_vec())),
3263 (b"key2".to_vec(), Some(b"value2".to_vec())),
3264 (b"key1".to_vec(), Some(b"value1".to_vec())),
3265 (b"key5".to_vec(), None),
3266 (b"key3".to_vec(), None),
3267 (b"key2".to_vec(), None),
3268 (b"key4".to_vec(), None),
3269 (b"key1".to_vec(), None),
3270 ]);
3271 test_basic(&[
3272 ([5u8; 250].to_vec(), Some(b"value5".to_vec())),
3273 ([5u8; 200].to_vec(), Some(b"value3".to_vec())),
3274 ([5u8; 100].to_vec(), Some(b"value4".to_vec())),
3275 ([5u8; 150].to_vec(), Some(b"value2".to_vec())),
3276 ([5u8; 101].to_vec(), Some(b"value1".to_vec())),
3277 ([5u8; 250].to_vec(), None),
3278 ([5u8; 101].to_vec(), None),
3279 ]);
3280 }
3281
3282 #[test]
3283 fn test_btree_iter() {
3284 let col_nb = 0;
3285 let mut data_start = Vec::new();
3286 for i in 0u8..100 {
3287 let mut key = b"key0".to_vec();
3288 key[3] = i;
3289 let mut value = b"val0".to_vec();
3290 value[3] = i;
3291 data_start.push((col_nb, key, Some(value)));
3292 }
3293 let mut data_change = Vec::new();
3294 for i in 0u8..100 {
3295 let mut key = b"key0".to_vec();
3296 if i % 2 == 0 {
3297 key[2] = i;
3298 let mut value = b"val0".to_vec();
3299 value[2] = i;
3300 data_change.push((col_nb, key, Some(value)));
3301 } else if i % 3 == 0 {
3302 key[3] = i;
3303 data_change.push((col_nb, key, None));
3304 } else {
3305 key[3] = i;
3306 let mut value = b"val0".to_vec();
3307 value[2] = i;
3308 data_change.push((col_nb, key, Some(value)));
3309 }
3310 }
3311
3312 let start_state: BTreeMap<Vec<u8>, Vec<u8>> =
3313 data_start.iter().cloned().map(|(_c, k, v)| (k, v.unwrap())).collect();
3314 let mut end_state = start_state.clone();
3315 for (_c, k, v) in data_change.iter() {
3316 if let Some(v) = v {
3317 end_state.insert(k.clone(), v.clone());
3318 } else {
3319 end_state.remove(k);
3320 }
3321 }
3322
3323 for stage in [
3324 EnableCommitPipelineStages::CommitOverlay,
3325 EnableCommitPipelineStages::LogOverlay,
3326 EnableCommitPipelineStages::DbFile,
3327 EnableCommitPipelineStages::Standard,
3328 ] {
3329 for i in 0..10 {
3330 test_btree_iter_inner(
3331 stage,
3332 &data_start,
3333 &data_change,
3334 &start_state,
3335 &end_state,
3336 i * 5,
3337 );
3338 }
3339 let data_start = vec![
3340 (0, b"key1".to_vec(), Some(b"val1".to_vec())),
3341 (0, b"key3".to_vec(), Some(b"val3".to_vec())),
3342 ];
3343 let data_change = vec![(0, b"key2".to_vec(), Some(b"val2".to_vec()))];
3344 let start_state: BTreeMap<Vec<u8>, Vec<u8>> =
3345 data_start.iter().cloned().map(|(_c, k, v)| (k, v.unwrap())).collect();
3346 let mut end_state = start_state.clone();
3347 for (_c, k, v) in data_change.iter() {
3348 if let Some(v) = v {
3349 end_state.insert(k.clone(), v.clone());
3350 } else {
3351 end_state.remove(k);
3352 }
3353 }
3354 test_btree_iter_inner(stage, &data_start, &data_change, &start_state, &end_state, 1);
3355 }
3356 }
3357 fn test_btree_iter_inner(
3358 db_test: EnableCommitPipelineStages,
3359 data_start: &[(u8, Vec<u8>, Option<Value>)],
3360 data_change: &[(u8, Vec<u8>, Option<Value>)],
3361 start_state: &BTreeMap<Vec<u8>, Vec<u8>>,
3362 end_state: &BTreeMap<Vec<u8>, Vec<u8>>,
3363 commit_at: usize,
3364 ) {
3365 let tmp = tempdir().unwrap();
3366 let mut options = db_test.options(tmp.path(), 5);
3367 let col_nb = 0;
3368 options.columns[col_nb as usize].btree_index = true;
3369 let db = Db::open_inner(&options, OpeningMode::Create).unwrap();
3370
3371 db.commit(data_start.iter().cloned()).unwrap();
3372 db_test.run_stages(&db);
3373
3374 let mut iter = db.iter(col_nb).unwrap();
3375 let mut iter_state = start_state.iter();
3376 let mut last_key = Value::new();
3377 for _ in 0..commit_at {
3378 let next = iter.next().unwrap();
3379 if let Some((k, _)) = next.as_ref() {
3380 last_key = k.clone();
3381 }
3382 assert_eq!(iter_state.next(), next.as_ref().map(|(k, v)| (k, v)));
3383 }
3384
3385 db.commit(data_change.iter().cloned()).unwrap();
3386 db_test.run_stages(&db);
3387
3388 let mut iter_state = end_state.range(last_key.clone()..);
3389 for _ in commit_at..100 {
3390 let mut state_next = iter_state.next();
3391 if let Some((k, _v)) = state_next.as_ref() {
3392 if *k == &last_key {
3393 state_next = iter_state.next();
3394 }
3395 }
3396 let iter_next = iter.next().unwrap();
3397 assert_eq!(state_next, iter_next.as_ref().map(|(k, v)| (k, v)));
3398 }
3399 let mut iter_state_rev = end_state.iter().rev();
3400 let mut iter = db.iter(col_nb).unwrap();
3401 iter.seek_to_last().unwrap();
3402 for _ in 0..100 {
3403 let next = iter.prev().unwrap();
3404 assert_eq!(iter_state_rev.next(), next.as_ref().map(|(k, v)| (k, v)));
3405 }
3406 }
3407
3408 #[cfg(feature = "instrumentation")]
3409 #[test]
3410 fn test_recover_from_log_on_error() {
3411 let tmp = tempdir().unwrap();
3412 let mut options = Options::with_columns(tmp.path(), 1);
3413 options.always_flush = true;
3414 options.with_background_thread = false;
3415
3416 {
3418 let db = Db::open_or_create(&options).unwrap();
3419 db.commit::<_, Vec<u8>>(vec![(0, vec![0], Some(vec![0]))]).unwrap();
3420 db.process_commits().unwrap();
3421 db.flush_logs().unwrap();
3422 db.enact_logs().unwrap();
3423 db.commit::<_, Vec<u8>>(vec![(0, vec![1], Some(vec![1]))]).unwrap();
3424 db.process_commits().unwrap();
3425 db.flush_logs().unwrap();
3426 crate::set_number_of_allowed_io_operations(4);
3427
3428 let err = db.enact_logs();
3430 assert!(err.is_err());
3431 db.inner.store_err(err);
3432 crate::set_number_of_allowed_io_operations(usize::MAX);
3433 }
3434
3435 {
3437 let db = Db::open(&options).unwrap();
3438 assert_eq!(db.get(0, &[0]).unwrap(), Some(vec![0]));
3439 assert_eq!(db.get(0, &[1]).unwrap(), Some(vec![1]));
3440 }
3441 }
3442
3443 #[cfg(feature = "instrumentation")]
3444 #[test]
3445 fn test_partial_log_recovery() {
3446 let tmp = tempdir().unwrap();
3447 let mut options = Options::with_columns(tmp.path(), 1);
3448 options.columns[0].btree_index = true;
3449 options.always_flush = true;
3450 options.with_background_thread = false;
3451
3452 {
3454 let db = Db::open_or_create(&options).unwrap();
3455 db.commit::<_, Vec<u8>>(vec![(0, vec![0], Some(vec![0]))]).unwrap();
3456 db.process_commits().unwrap();
3457 db.commit::<_, Vec<u8>>(vec![(0, vec![1], Some(vec![1]))]).unwrap();
3458 crate::set_number_of_allowed_io_operations(4);
3459 assert!(db.process_commits().is_err());
3460 crate::set_number_of_allowed_io_operations(usize::MAX);
3461 db.flush_logs().unwrap();
3462 }
3463
3464 {
3466 let db = Db::open(&options).unwrap();
3467 assert_eq!(db.get(0, &[0]).unwrap(), Some(vec![0]));
3468 }
3469
3470 {
3472 let db = Db::open(&options).unwrap();
3473 assert!(db.get(0, &[0]).unwrap().is_some());
3474 }
3475 }
3476
3477 #[cfg(feature = "instrumentation")]
3478 #[test]
3479 fn test_continue_reindex() {
3480 let _ = env_logger::try_init();
3481 let tmp = tempdir().unwrap();
3482 let mut options = Options::with_columns(tmp.path(), 1);
3483 options.columns[0].preimage = true;
3484 options.columns[0].uniform = true;
3485 options.always_flush = true;
3486 options.with_background_thread = false;
3487 options.salt = Some(Default::default());
3488
3489 {
3490 let db = Db::open_or_create(&options).unwrap();
3492 let commit: Vec<_> = (0..65u32)
3493 .map(|index| {
3494 let mut key = [0u8; 32];
3495 key[2] = (index as u8) << 1;
3496 (0, key.to_vec(), Some(vec![index as u8]))
3497 })
3498 .collect();
3499 db.commit(commit).unwrap();
3500
3501 db.process_commits().unwrap();
3502 db.flush_logs().unwrap();
3503 db.enact_logs().unwrap();
3504 std::fs::copy(tmp.path().join("index_00_16"), tmp.path().join("index_00_16.bak"))
3509 .unwrap();
3510 db.process_reindex().unwrap();
3511 db.flush_logs().unwrap();
3512 db.enact_logs().unwrap();
3513 db.clean_logs().unwrap();
3514 std::fs::rename(tmp.path().join("index_00_16.bak"), tmp.path().join("index_00_16"))
3515 .unwrap();
3516 }
3517
3518 {
3520 let db = Db::open(&options).unwrap();
3521 db.process_reindex().unwrap();
3522 let mut entries = 0;
3523 db.iter_column_while(0, |_| {
3524 entries += 1;
3525 true
3526 })
3527 .unwrap();
3528
3529 assert_eq!(entries, 65);
3530 assert_eq!(db.inner.columns[0].index_bits(), Some(17));
3531 }
3532 }
3533
3534 #[test]
3535 fn test_remove_column() {
3536 let tmp = tempdir().unwrap();
3537 let db_test_file = EnableCommitPipelineStages::DbFile;
3538 let mut options_db_files = db_test_file.options(tmp.path(), 2);
3539 options_db_files.salt = Some(options_db_files.salt.unwrap_or_default());
3540 let mut options_std = EnableCommitPipelineStages::Standard.options(tmp.path(), 2);
3541 options_std.salt = options_db_files.salt.clone();
3542
3543 let db = Db::open_inner(&options_db_files, OpeningMode::Create).unwrap();
3544
3545 let payload: Vec<(u8, _, _)> = (0u16..100)
3546 .map(|i| (1, i.to_le_bytes().to_vec(), Some(i.to_be_bytes().to_vec())))
3547 .collect();
3548
3549 db.commit(payload.clone()).unwrap();
3550
3551 db_test_file.run_stages(&db);
3552 drop(db);
3553
3554 let db = Db::open_inner(&options_std, OpeningMode::Write).unwrap();
3555 for (col, key, value) in payload.iter() {
3556 assert_eq!(db.get(*col, key).unwrap().as_ref(), value.as_ref());
3557 }
3558 drop(db);
3559 Db::reset_column(&mut options_db_files, 1, None).unwrap();
3560
3561 let db = Db::open_inner(&options_db_files, OpeningMode::Write).unwrap();
3562 for (col, key, _value) in payload.iter() {
3563 assert_eq!(db.get(*col, key).unwrap(), None);
3564 }
3565
3566 let payload: Vec<(u8, _, _)> = (0u16..10)
3567 .map(|i| (1, i.to_le_bytes().to_vec(), Some(i.to_be_bytes().to_vec())))
3568 .collect();
3569
3570 db.commit(payload.clone()).unwrap();
3571
3572 db_test_file.run_stages(&db);
3573 drop(db);
3574
3575 let db = Db::open_inner(&options_std, OpeningMode::Write).unwrap();
3576 let payload: Vec<(u8, _, _)> = (10u16..100)
3577 .map(|i| (1, i.to_le_bytes().to_vec(), Some(i.to_be_bytes().to_vec())))
3578 .collect();
3579
3580 db.commit(payload.clone()).unwrap();
3581 assert!(db.iter(1).is_err());
3582
3583 drop(db);
3584
3585 let mut col_option = options_std.columns[1].clone();
3586 col_option.btree_index = true;
3587 Db::reset_column(&mut options_std, 1, Some(col_option)).unwrap();
3588
3589 let db = Db::open_inner(&options_std, OpeningMode::Write).unwrap();
3590 let payload: Vec<(u8, _, _)> = (0u16..10)
3591 .map(|i| (1, i.to_le_bytes().to_vec(), Some(i.to_be_bytes().to_vec())))
3592 .collect();
3593
3594 db.commit(payload.clone()).unwrap();
3595 assert!(db.iter(1).is_ok());
3596 }
3597}