Skip to main content

parity_db/
db.rs

1// Copyright 2021-2022 Parity Technologies (UK) Ltd.
2// This file is dual-licensed as Apache-2.0 or MIT.
3
4//! The database objects is split into `Db` and `DbInner`.
5//! `Db` creates shared `DbInner` instance and manages background
6//! worker threads that all use the inner object.
7//!
8//! There are 4 worker threads:
9//! log_worker: Processes commit queue and reindexing. For each commit
10//! in the queue, log worker creates a write-ahead record using `Log`.
11//! Additionally, if there are active reindexing, it creates log records
12//! for batches of relocated index entries.
13//! flush_worker: Flushes log records to disk by calling `fsync` on the
14//! log files.
15//! commit_worker: Reads flushed log records and applies operations to the
16//! index and value tables.
17//! cleanup_worker: Flush tables by calling `fsync`, and cleanup log.
18//! Each background worker is signalled with a conditional variable once
19//! there is some work to be done.
20
21use 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
53// Max size of commit queue. (Keys + Values). If the queue is
54// full `commit` will block.
55// These are in memory, so we use usize
56const MAX_COMMIT_QUEUE_BYTES: usize = 16 * 1024 * 1024;
57// Max size of log overlay. If the overlay is full, processing
58// of commit queue is blocked.
59const MAX_LOG_QUEUE_BYTES: i64 = 128 * 1024 * 1024;
60// Minimum size of log file before it is considered full.
61const MIN_LOG_SIZE_BYTES: u64 = 64 * 1024 * 1024;
62// Number of log files to keep after flush when sync mode is disabled. Give the database some chance
63// to recover in case of crash.
64const KEEP_LOGS: usize = 16;
65// Hard limit on the number of log files in sync mode. The number of log may grow while existing
66// logs are waiting on fsync. Commits will be throttled if total number of log files exceeds this
67// number.
68const MAX_LOG_FILES: usize = 4;
69
70/// Value is just a vector of bytes. Value sizes up to 4Gb are allowed.
71pub 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// Commit data passed to `commit`
157#[derive(Debug, Default)]
158struct Commit {
159	// Commit ID. This is not the same as log record id, as some records
160	// are originated within the DB. E.g. reindex.
161	id: u64,
162	// Size of user data pending insertion (keys + values) or
163	// removal (keys)
164	bytes: usize,
165	// Operations.
166	changeset: CommitChangeSet,
167}
168
169// Pending commits. This may not grow beyond `MAX_COMMIT_QUEUE_BYTES` bytes.
170#[derive(Debug, Default)]
171struct CommitQueue {
172	// Log record.
173	record_id: u64,
174	// Total size of all commits in the queue.
175	bytes: usize,
176	// FIFO queue.
177	commits: VecDeque<Commit>,
178}
179
180#[derive(Debug)]
181struct Trees {
182	readers: HashMap<Key, Weak<RwLock<Box<dyn TreeReader + Send + Sync>>>, IdentityBuildHasher>,
183	/// Number of queued dereferences for each tree
184	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	// Overlay of most recent values in the commit queue.
198	commit_overlay: RwLock<Vec<CommitOverlay>>,
199	trees: RwLock<HashMap<ColId, Trees>>,
200	// This may underflow occasionally, but is bound for 0 eventually.
201	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				// Check commit overlay first
316				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				// Go into tables and log overlay.
321				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				// We lock log, if btree structure changed while reading that would be an issue.
331				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				// Check commit overlay first
348				if let Some(l) = overlay.get(col as usize).and_then(|o| o.get_size(&key)) {
349					return Ok(l)
350				}
351				// Go into tables and log overlay.
352				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				// Check commit overlay first
406				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				// Check commit overlay first
442				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				// Check if the tree actually exists. We can't return the data from this function as
470				// TreeReader is not locked. That is done by the client.
471				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	// Commit simply adds the data to the queue and to the overlay and
522	// exits as early as possible.
523	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										// Check if TreeReader is active for this tree
571										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				// Cleanup the commit overlay with old id.
775				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			// Wait if the queue is full.
808			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					// Past the waiting threshold.
828					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							// Check if there are currently any locks on the tree. Will need to
847							// defer if there are.
848							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							// Also check if there are any later commits in the queue that use this
867							// tree. Will need to defer if there are.
868							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						// Generate a new id
886						None
887					} else {
888						// Nothing else in the queue so can reuse same id
889						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			// Collect final changes to value tables
948			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				// Cleanup the commit overlay.
964				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		// Process any pending reindexes
1002		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					// Validate all records before applying anything
1130					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									// Check if there's another reindex on the next iteration
1230									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									// Check if there's another reindex on the next iteration
1245									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		// Re-read any cached metadata
1352		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				// On error the log reader may be left in inconsistent state. So it is important
1372				// to no attempt any further log enactment.
1373				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		// Finish logged records and proceed to log and enact queued commits.
1380		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
1459/// Database instance.
1460pub 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	/// Open the database with given options. An error will be returned if the database does not
1476	/// exist.
1477	pub fn open(options: &Options) -> Result<Db> {
1478		Self::open_inner(options, OpeningMode::Write)
1479	}
1480
1481	/// Open the database using given options. If the database does not exist it will be created
1482	/// empty.
1483	pub fn open_or_create(options: &Options) -> Result<Db> {
1484		Self::open_inner(options, OpeningMode::Create)
1485	}
1486
1487	/// Open an existing database in read-only mode.
1488	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		// This needs to be call before log thread: so first reindexing
1496		// will run in correct state.
1497		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	/// Get a value in a specified column by key. Returns `None` if the key does not exist.
1551	pub fn get(&self, col: ColId, key: &[u8]) -> Result<Option<Value>> {
1552		self.inner.get(col, key, true)
1553	}
1554
1555	/// Get value size by key. Returns `None` if the key does not exist.
1556	pub fn get_size(&self, col: ColId, key: &[u8]) -> Result<Option<u32>> {
1557		self.inner.get_size(col, key)
1558	}
1559
1560	/// Iterate over all ordered key-value pairs. Only supported for columns configured with
1561	/// `btree_indexed`.
1562	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	/// Commit a set of changes to the database.
1595	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	/// Commit a set of changes to the database.
1604	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	/// Commit a set of changes to the database.
1612	///
1613	/// This method passes values as `Arc<Vec<u8>>` potentially eliminating an extra copy.
1614	#[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	/// Commit a set of changes to the database.
1626	///
1627	/// This method passes values as `Bytes` potentially eliminating an extra copy.
1628	#[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	/// Returns the number of columns in the database.
1641	pub fn num_columns(&self) -> u8 {
1642		self.inner.columns.len() as u8
1643	}
1644
1645	/// Iterate a column and call a function for each value. This is only supported for columns with
1646	/// `btree_index` set to `false`. Iteration order is unspecified.
1647	/// Unlike `get` the iteration may not include changes made in recent `commit` calls.
1648	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	/// Iterate a column and call a function for each value. This is only supported for columns with
1653	/// `btree_index` set to `false`. Iteration order is unspecified. Note that the
1654	/// `key` field in the state is the hash of the original key.
1655	/// Unlike `get` the iteration may not include changes made in recent `commit` calls.
1656	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		// Start with pending reindex.
1682		let mut more_reindex = db.process_reindex()?;
1683		let mut more_commits = false;
1684		// Process all commits but allow reindex to be interrupted.
1685		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	/// Dump full database stats to the text output.
1722	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	/// Reset internal database statistics for the database or specified column.
1731	pub fn clear_stats(&self, column: Option<u8>) -> Result<()> {
1732		self.inner.clear_stats(column)
1733	}
1734
1735	/// Print database contents in text form to stderr.
1736	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	/// Get database statistics.
1748	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	// We open the DB before to check metadata validity and make sure there are no pending WAL
1764	// logs.
1765	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	/// Add a new column with options specified by `new_column_options`.
1773	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	/// Remove last column from the database.
1783	/// Db must be close when called.
1784	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	/// Truncate a column from the database, optionally changing its options.
1798	/// Db must be close when called.
1799	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
1896// Use a trait here to allow client code to have better control over lock guard lifetime without
1897// lifetime proliferation within Db (which would be required if not using a dynamic object).
1898pub 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		/* let value = self.db.get(self.col, &self.key)?;
1914		if let Some(data) = value {
1915			return unpack_node_data(data).map(|x| Some(x))
1916		}
1917		Err(Error::InvalidValueData) */
1918
1919		match &self.db.columns[self.col as usize] {
1920			Column::Hash(column) => {
1921				let overlay = self.db.commit_overlay.read();
1922				// Check commit overlay first
1923				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					// Go into tables and log overlay.
1929					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/// Different operations allowed for a commit.
2049/// Behavior may differs depending on column configuration.
2050#[derive(Debug, PartialEq, Eq)]
2051pub enum Operation<Key, Value> {
2052	/// Insert or update the value for a given key.
2053	Set(Key, Value),
2054
2055	/// Dereference at a given key, resulting in
2056	/// either removal of a key value or decrement of its
2057	/// reference count counter.
2058	Dereference(Key),
2059
2060	/// Increment the reference count counter of an existing value for a given key.
2061	/// If no value exists for the key, this operation is skipped.
2062	Reference(Key),
2063
2064	/// Insert a new tree into a MultiTree column using root key and node structure.
2065	InsertTree(Key, NewNode),
2066
2067	/// Increment the reference count of a tree (at root Key) from a MultiTree column.
2068	ReferenceTree(Key),
2069
2070	/// Dereference an existing tree (at root Key) from a MultiTree column, resulting in either
2071	/// removal of the tree or decrement of its reference count.
2072	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	/// (address, value)
2127	NewValue(u64, RcValue),
2128	/// (address)
2129	IncrementReference(u64),
2130	/// Dereference and remove any of the children in the tree
2131	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					// Don't add removed ref-counted values to overlay.
2211					if !ref_counted {
2212						overlay.indexed.insert(*k, (record_id, None));
2213					}
2214				},
2215				Operation::Reference(..) => {
2216					// Don't add (we allow remove value in overlay when using rc: some
2217					// indexing on top of it is expected).
2218					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 has triggered another reindex.
2259				*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					// TODO: Remove TreeReader from Db.
2302				},
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			// Can't move this after write_address_dec_ref_plan as write_address_dec_ref_plan might
2318			// free the node meaning it could get reclaimed. Then get_node_children will return
2319			// incorrect data.
2320			let node = guard.get_node_children(*address)?;
2321			let (remains, _outcome) = column.write_address_dec_ref_plan(*address, writer)?;
2322			if !remains {
2323				// Was removed
2324				*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
2370/// Verification operation utilities.
2371pub mod check {
2372	/// Database dump verbosity.
2373	pub enum CheckDisplay {
2374		/// Don't output any data.
2375		None,
2376		/// Output full data.
2377		Full,
2378		/// Limit value output to the specified size.
2379		Short(u64),
2380	}
2381
2382	/// Options for producing a database dump.
2383	pub struct CheckOptions {
2384		/// Only process this column. If this is `None` all columns will be processed.
2385		pub column: Option<u8>,
2386		/// Start with this index.
2387		pub from: Option<u64>,
2388		/// End with this index.
2389		pub bound: Option<u64>,
2390		/// Verbosity.
2391		pub display: CheckDisplay,
2392		/// Ordered validation.
2393		pub fast: bool,
2394		/// Make sure free lists are correct.
2395		pub validate_free_refs: bool,
2396	}
2397
2398	impl CheckOptions {
2399		/// Create a new instance.
2400		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	// This is used in tests to disable certain commit stages.
2445	#[derive(Eq, PartialEq, Debug, Clone, Copy)]
2446	enum EnableCommitPipelineStages {
2447		// No threads started, data stays in commit overlay.
2448		#[allow(dead_code)]
2449		CommitOverlay,
2450		// Log worker run, data processed up to the log overlay.
2451		#[allow(dead_code)]
2452		LogOverlay,
2453		// Runing all.
2454		#[allow(dead_code)]
2455		DbFile,
2456		// Default run mode.
2457		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									// the signal is triggered just before cleaning the overlay, so
2500									// we wait a bit.
2501									// Warning this is still rather flaky and should be ignored
2502									// or removed.
2503									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		// issue with some file reopening when no delay
2643		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		// Expected number of columns
2703		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		// Write to new columns.
2709		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		// Reopen DB and fetch all keys we inserted.
2721		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		// Fetch from new columns
2728		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		// issue with some file reopening when no delay
2859		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		// ref counted and commit overlay currently don't support removal
3010		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			// btree do not have stats implemented
3063			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(); // should be out of loop, yet it makes alternance insert/delete in some case.
3091			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		// We do 2 commits and we fail while enacting the second one
3417		{
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			// Set the background error explicitly as background threads are disabled in tests.
3429			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		// Open the databases and check that both values are there.
3436		{
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		// We do 2 commits and we fail while writing the second one
3453		{
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		// We open a first time, the first value is there
3465		{
3466			let db = Db::open(&options).unwrap();
3467			assert_eq!(db.get(0, &[0]).unwrap(), Some(vec![0]));
3468		}
3469
3470		// We open a second time, the first value should be still there
3471		{
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			// Force a reindex by committing more than 64 values with the same 16 bit prefix
3491			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			// i16 now contains 64 values and i17 contains a single value that did not fit
3505
3506			// Simulate interrupted reindex by processing it first and then restoring the old index
3507			// file. Make a copy of the index file first.
3508			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		// Reopen the database which should load the reindex.
3519		{
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}