diff --git a/Cargo.toml b/Cargo.toml index da763fc..f9e9146 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -8,8 +8,14 @@ tempdir = "0.3.7" tempfile = "3.14.0" criterion = "0.5.1" + [dependencies] chunkfs = { version = "0.1", features = ["chunkers", "hashers"] } gnuplot = "0.0.44" approx = "0.5.1" rand = "0.8.5" +tokio = { version = "1.44.2", features = ["full"] } +serde = { version = "1.0", features = ["derive"] } +bincode = "1.3" +async-recursion = "1.1.1" +futures = "0.3.31" diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index d5c1652..6bfe1ed 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -1,21 +1,155 @@ -use chunkfs::{Data, DataContainer, Database}; - use std::{ - cell::RefCell, - fmt::{self, Debug}, + collections::{HashSet, VecDeque}, + fmt::Debug, fs::{create_dir_all, File}, - io::{self, ErrorKind}, + io::{self, BufReader, BufWriter, ErrorKind}, + mem, os::unix::fs::FileExt, - path::PathBuf, + path::{Path, PathBuf}, rc::Rc, + sync::{ + atomic::{AtomicU64, AtomicUsize, Ordering}, + Arc, Mutex, + }, + thread, time, }; +use async_recursion::async_recursion; + +use serde::{Deserialize, Serialize}; + +use chunkfs::{Data, DataContainer, Database}; +use tokio::{self, runtime::Runtime, sync::RwLock}; + const DEFAULT_MAX_FILE_SIZE: u64 = 2 << 20; +pub trait BPlusKey: Default + Ord + Clone + Sized + Sync + Send {} +impl BPlusKey for T {} + +pub trait BPlusKeySerializable: BPlusKey + Serialize + for<'de> Deserialize<'de> {} +impl Deserialize<'de>> + BPlusKeySerializable for T +{ +} + extern crate chunkfs; +/// Serializable version of BPlusTree +#[derive(Serialize, Deserialize)] +struct SerializableBPlus { + t: usize, + path: PathBuf, + file_number: usize, + offset: u64, + max_file_size: u64, + root: SerializableNode, +} + +/// Easily serializable version of BPlusTree Node +#[derive(Serialize, Deserialize)] +enum SerializableNode { + Internal(SerializableInternalNode), + Leaf(SerializableLeaf), +} + +#[derive(Serialize, Deserialize)] +struct SerializableInternalNode { + keys: Vec, + children: Vec>, +} + +#[derive(Serialize, Deserialize)] +struct SerializableLeaf { + entries: Vec<(K, ChunkHandler)>, +} + +impl BPlus { + /// Returns new instance of SerializableBPlus with data from provided BPlus + async fn serialize(&self) -> SerializableBPlus { + SerializableBPlus { + t: self.t, + path: self.path.clone(), + file_number: self.file_number.load(Ordering::SeqCst), + offset: self.offset.load(Ordering::SeqCst), + max_file_size: self.max_file_size, + root: self.root.read().await.serialize().await, + } + } +} + +impl Node { + #[async_recursion] + /// Returns new instance of SerializableNode with data from provided Node + async fn serialize(&self) -> SerializableNode { + match self { + Node::Internal(internal) => { + let keys = internal.keys.iter().map(|k| (**k).clone()).collect(); + + let children_clone = internal.children.clone(); + let mut children = Vec::new(); + for child in children_clone { + children.push(child.read().await.serialize().await); + } + + SerializableNode::Internal(SerializableInternalNode { keys, children }) + } + Node::Leaf(leaf) => SerializableNode::Leaf(SerializableLeaf { + entries: leaf + .entries + .iter() + .map(|(k, v)| ((**k).clone(), v.clone())) + .collect(), + }), + } + } +} + +impl SerializableBPlus { + /// Returns new instance of BPlus with data from provided BPlusSerializable + async fn deserialize(self) -> BPlus { + let root = Arc::new(RwLock::new(Node::from(self.root))); + + let tree = BPlus { + root: root.clone(), + t: self.t, + path: self.path.clone(), + file_number: AtomicUsize::new(self.file_number), + offset: AtomicU64::new(self.offset), + current_file: BPlus::::open_current_file(&self.path, self.file_number).unwrap(), + max_file_size: self.max_file_size, + latch: RwLock::new(()), + }; + + tree.rebuild_links().await; + tree + } +} + +impl From> for Node { + fn from(node: SerializableNode) -> Self { + match node { + SerializableNode::Internal(internal) => Node::Internal(InternalNode { + keys: internal.keys.into_iter().map(Arc::new).collect(), + children: internal + .children + .into_iter() + .map(|c| Arc::new(RwLock::new(Node::from(c)))) + .collect(), + }), + SerializableNode::Leaf(leaf) => Node::Leaf(Leaf { + entries: leaf + .entries + .into_iter() + .map(|(k, v)| (Arc::new(k), v)) + .collect(), + next: None, + }), + } + } +} + /// Structure that handles chunks written in files. -#[derive(Clone, Default, Debug)] +#[derive(Clone, Default, Debug, Serialize, Deserialize)] pub struct ChunkHandler { /// Path to file with chunk. path: PathBuf, @@ -43,7 +177,7 @@ impl ChunkHandler { } /// A type that represents a reference to another node. -type Link = Rc>>; +type Link = Arc>>; /// Represents a node in a B+ tree. /// All data resides in leaf nodes, while internal nodes. @@ -60,14 +194,14 @@ struct InternalNode { /// Children of that node. children: Vec>, /// Keys of that node. - keys: Vec>, + keys: Vec>, } /// Leaf node in a B+ tree #[derive(Default, Clone)] struct Leaf { /// Data entries that stored in that leaf. - entries: Vec<(Rc, ChunkHandler)>, + entries: Vec<(Arc, ChunkHandler)>, /// Link to the next leaf; None if there are none. next: Option>, } @@ -75,121 +209,288 @@ struct Leaf { /// B+ tree pub struct BPlus { /// Root of the B+ tree. - root: Node, + root: Link, /// Parameter, that represents minimal and maximal amount of node keys. t: usize, /// Path to the directory, in which all data will be writen. path: PathBuf, /// Number of current file. - file_number: usize, + file_number: AtomicUsize, /// Current offset in current file. - offset: u64, + offset: AtomicU64, /// Current file. - current_file: File, + current_file: Arc>, /// Max file size. max_file_size: u64, + // Latch for root + latch: RwLock<()>, +} + +/// Wrapper for BPlusTree with sync functions with async runtime +pub struct BPlusStorage { + /// BPlusTree + tree: Arc>, + /// Async tokio runtime for operations + runtime: Runtime, + /// Currently inserting keys + keys_set: Arc>>, } -impl fmt::Display for BPlus { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - BPlus::fmt_node(&self.root, 0, f) +impl BPlusStorage { + /// Creates new instance of B+ tree with given runtime, t and path + /// + /// runtime is tokio runtime + /// + /// t represents minimal and maximum quantity of keys in the node + /// + /// All data will be written in directory by given path + pub fn new(runtime: Runtime, t: usize, path: PathBuf) -> io::Result { + let tree = BPlus::new(t, path).unwrap(); + Ok(Self { + tree: Arc::new(tree), + runtime, + keys_set: Arc::new(Mutex::new(HashSet::new())), + }) } } -impl BPlus { - fn fmt_node(node: &Node, level: usize, f: &mut fmt::Formatter<'_>) -> fmt::Result { - match node { - Node::Internal(internal) => { - writeln!( - f, - "{}[Internal] keys: {:?}", - " ".repeat(level), - internal.keys - )?; - - for child in &internal.children { - BPlus::fmt_node(&child.borrow(), level + 1, f)?; +impl Database> for BPlusStorage { + /// Inserts given value by given key in the B+ tree + fn insert(&mut self, key: K, value: DataContainer<()>) -> io::Result<()> { + let tree = self.tree.clone(); + + let value = match value.extract() { + Data::Chunk(chunk) => chunk.clone(), + Data::TargetChunk(_chunk) => unimplemented!(), + }; + + let set_clone = self.keys_set.clone(); + set_clone.lock().unwrap().insert(key.clone()); + + self.runtime.spawn(async move { + tree.insert(key.clone(), value).await; + set_clone.lock().unwrap().remove(&key); + }); + Ok(()) + } + + /// Gets value by given key from B+ tree + fn get(&self, key: &K) -> io::Result> { + let tree = self.tree.clone(); + let set_clone = self.keys_set.clone(); + + Ok(self + .runtime + .block_on(async move { + while set_clone.lock().unwrap().contains(key) { + thread::sleep(time::Duration::from_millis(10)); } - Ok(()) - } - Node::Leaf(leaf) => { - writeln!( - f, - "{}[Leaf] entries: {:?}", - " ".repeat(level), - leaf.entries - ) - } - } + tree.get(key).await.unwrap() + }) + .into()) + } + + /// Returns whether key is contained in the B+ tree or not + fn contains(&self, key: &K) -> bool { + self.get(key).is_ok() } } #[allow(dead_code)] -impl BPlus { +impl BPlus { /// Creates new instance of B+ tree with given t and path + /// /// t represents minimal and maximal quantity of keys in node + /// /// All data will be written in files in directory by given path pub fn new(t: usize, path: PathBuf) -> io::Result { let path_to_file = path.join("0"); create_dir_all(&path)?; let current_file = File::create(path_to_file)?; + Ok(Self { - root: Node::Leaf(Leaf::default()), + root: Arc::new(RwLock::new(Node::Leaf(Leaf::default()))), t, path, - file_number: 0, - offset: 0, - current_file, + file_number: 0.into(), + offset: 0.into(), + current_file: Arc::new(RwLock::new(current_file)), max_file_size: DEFAULT_MAX_FILE_SIZE, + latch: RwLock::new(()), }) } /// Creates new chunk_handler and writes data to a file - fn get_chunk_handler(&mut self, value: Vec) -> io::Result { - if self.offset >= self.max_file_size { - self.file_number += 1; - self.offset = 0; - self.current_file = File::create(self.path.join(self.file_number.to_string())).unwrap(); + async fn get_chunk_handler(&self, value: Vec) -> io::Result { + let mut file_guard = self.current_file.write().await; + if self.offset.load(std::sync::atomic::Ordering::SeqCst) >= self.max_file_size { + self.file_number + .fetch_add(1, std::sync::atomic::Ordering::SeqCst); + self.offset.store(0, std::sync::atomic::Ordering::SeqCst); + let file_number = self.file_number.load(Ordering::SeqCst).to_string(); + let file_path = self.path.join(file_number); + + *file_guard = File::create(file_path).unwrap(); } let value_size = value.len(); - self.current_file.write_at(&value, self.offset)?; + file_guard.write_at( + &value, + self.offset.load(std::sync::atomic::Ordering::SeqCst), + )?; let value_to_insert = ChunkHandler::new( - self.path.join(self.file_number.to_string()), - self.offset, + self.path.join( + self.file_number + .load(std::sync::atomic::Ordering::SeqCst) + .to_string(), + ), + self.offset.load(std::sync::atomic::Ordering::SeqCst), value.len(), ); - self.offset += value_size as u64; + self.offset + .fetch_add(value_size as u64, std::sync::atomic::Ordering::SeqCst); Ok(value_to_insert) } /// Inserts given value by given key in the B+ tree /// /// Returns Err(_) if file could not be created - pub fn insert(&mut self, key: K, value: Vec) -> io::Result<()> { - let value_to_insert = self.get_chunk_handler(value).unwrap(); + pub async fn insert(&self, key: K, value: Vec) { + let value = self.get_chunk_handler(value).await.unwrap(); + let mut path = Vec::new(); // Path to leaf + // Insert that implies that target leaf is safe. Otherwise returns Err() + if self + .optimistic_insert(key.clone(), value.clone()) + .await + .is_ok() + { + return; + } + let mut latch_guard = Some(self.latch.write()); + let key = Arc::new(key); + let mut current = self.root.clone(); + let mut split_result; + let mut guards = VecDeque::new(); - let key = Rc::new(key); + // Descent to the leaf + loop { + let mut current_node = current.write_owned().await; + if let Some(guard) = latch_guard.take() { + drop(guard); + latch_guard = None; + }; + match &mut *current_node { + Node::Leaf(leaf) => { + match leaf.entries.binary_search_by(|(k, _)| k.cmp(&key)) { + Ok(pos) => leaf.entries[pos] = (key.clone(), value), + Err(pos) => leaf.entries.insert(pos, (key.clone(), value)), + }; + + split_result = if leaf.entries.len() == 2 * self.t { + Some(current_node.split(self.t)) + } else { + while !guards.is_empty() { + drop(guards.pop_front().unwrap()); + } + None + }; + + // if path is empty, then current node is root + if path.is_empty() { + guards.push_back(current_node); + } else { + drop(current_node); + } - if let Some((new_node_link, new_key)) = self.root.insert(key, value_to_insert, self.t) { - let mut prev_root = self.root.clone(); - if let Node::Leaf(mut leaf) = prev_root { - leaf.next = Some(new_node_link.clone()); - prev_root = Node::Leaf(leaf); + break; + } + Node::Internal(internal) => { + let pos = match internal.keys.binary_search(&key) { + Ok(pos) => pos + 1, + Err(pos) => pos, + }; + + // droping guards if nodes are not going to be changed + if internal.keys.len() != 2 * self.t - 2 { + while !guards.is_empty() { + drop(guards.pop_front().unwrap()); + } + } + + let next_node = internal.children[pos].clone(); + + path.push(pos); + + current = next_node; + } } - let mut new_root_children = Vec::with_capacity(2 * self.t); - let mut new_root_keys = Vec::with_capacity(2 * self.t - 1); - new_root_children.push(Rc::new(RefCell::new(prev_root))); - new_root_children.push(new_node_link); - new_root_keys.push(new_key); - let new_root = Node::::Internal(InternalNode { - children: new_root_children, - keys: new_root_keys, - }); - - self.root = new_root; + + guards.push_back(current_node); + } + + // Going up to the root splitting nodes if needed + while let Some(pos) = path.pop() { + if let Some((new_node, median)) = split_result.take() { + let mut node = guards.pop_back().unwrap(); + if let Node::Internal(internal) = &mut *node { + internal.keys.insert(pos, median.clone()); + internal.children.insert(pos + 1, new_node); + if internal.keys.len() == 2 * self.t - 1 { + split_result = Some(node.split(self.t)); + } else { + split_result = None; + } + } + if path.is_empty() { + guards.push_back(node); + } else { + drop(node); + } + } + } + + // splitting root if needed + if let Some((new_node, median)) = split_result.take() { + // if path is empty, then current node is root + if path.is_empty() { + if let Some(mut node) = guards.pop_back() { + match &mut *node { + Node::Internal(internal) => { + let mut old_root_children = Vec::new(); + let mut old_root_keys = Vec::new(); + mem::swap(&mut old_root_keys, &mut internal.keys); + mem::swap(&mut old_root_children, &mut internal.children); + let old_root = Node::::Internal(InternalNode { + children: (old_root_children), + keys: (old_root_keys), + }); + internal.children.push(Arc::new(RwLock::new(old_root))); + internal.children.push(new_node); + internal.keys.push(median.clone()); + } + Node::Leaf(leaf) => { + let mut old_root_entries = Vec::new(); + let old_root_next = leaf.next.clone(); + mem::swap(&mut old_root_entries, &mut leaf.entries); + let old_root = Node::::Leaf(Leaf { + entries: old_root_entries, + next: old_root_next, + }); + let new_root = Node::::Internal(InternalNode { + children: (vec![Arc::new(RwLock::new(old_root)), new_node]), + keys: (vec![median.clone()]), + }); + *node = new_root; + } + } + drop(node); + } + } + } + + for guard in guards { + drop(guard); } - Ok(()) } #[allow(unused_variables)] @@ -198,304 +499,370 @@ impl BPlus { } /// Gets value from a B+ tree by given key - pub fn get(&self, key: &K) -> io::Result> { - self.root.get(key) + pub async fn get(&self, key: &K) -> io::Result> { + let mut latch_guard = Some(self.latch.read()); + let mut current = self.root.clone(); + + let mut prev_guard = None; + loop { + let node = current.read_owned().await; + if let Some(guard) = latch_guard { + drop(guard); + latch_guard = None; + } + if prev_guard.is_some() { + drop(prev_guard); + } + match &*node { + Node::Leaf(leaf) => { + return match leaf.entries.binary_search_by(|(k, _)| k.as_ref().cmp(key)) { + Ok(pos) => { + let data_read_result = leaf.entries[pos].1.read()?; + drop(node); + Ok(data_read_result) + } + Err(_) => { + drop(node); + Err(ErrorKind::NotFound.into()) + } + }; + } + Node::Internal(internal) => { + let pos = match internal.keys.binary_search_by(|k| k.as_ref().cmp(key)) { + Ok(pos) => pos + 1, + Err(pos) => pos, + }; + + current = match internal.children.get(pos) { + Some(child) => child.clone(), + None => { + drop(node); + return Err(ErrorKind::NotFound.into()); + } + }; + } + } + prev_guard = Some(node); + } } -} -impl Node { - /// Splits node into two and returns new node with it first key - fn split(&mut self, t: usize) -> (Link, Rc) { - match self { - Node::Leaf(leaf) => { - let mut new_leaf_entries = leaf.entries.split_off(t); - new_leaf_entries.reserve_exact(t); - let middle_key = new_leaf_entries[0].0.clone(); + /// For optimistic latch crabbing + /// + /// Insert firstly implies that leaf is safe + /// + /// If it is safe, than inserts(without write locks on other nodes) to the leaf and returns Ok + /// + /// Else, returns Err + /// + /// Also returns Err if root is leaf + async fn optimistic_insert(&self, key: K, value: ChunkHandler) -> Result<(), ()> { + let mut latch_guard = Some(self.latch.read()); + let mut current = self.root.clone(); + let key = Arc::new(key); - let new_leaf = Node::Leaf(Leaf { - entries: new_leaf_entries, - next: leaf.next.take(), - }); + let mut prev_guard = None; + let mut last_child_index = None; - let new_leaf_link = Rc::new(RefCell::new(new_leaf.clone())); - leaf.next = Some(new_leaf_link.clone()); + loop { + let node = current.read_owned().await; - (new_leaf_link, middle_key) + if let Some(guard) = latch_guard.take() { + drop(guard); + if matches!(&*node, Node::Leaf(_)) { + return Err(()); + } } - Node::Internal(internal_node) => { - let mut new_node_keys = internal_node.keys.split_off(t - 1); - let middle_key = new_node_keys.remove(0); - let mut new_node_children = internal_node.children.split_off(t); - new_node_keys.reserve_exact(t); - new_node_children.reserve_exact(t); + if matches!(&*node, Node::Leaf(_)) { + break; + } - let new_node = Node::Internal(InternalNode { - children: new_node_children, - keys: new_node_keys, - }); + prev_guard = Some(node); - (Rc::new(RefCell::new(new_node)), middle_key) + if let Node::Internal(internal) = prev_guard.as_deref().unwrap() { + let pos = match internal.keys.binary_search(&key) { + Ok(pos) => pos + 1, + Err(pos) => pos, + }; + last_child_index = Some(pos); + current = internal.children[pos].clone(); + } else { + unreachable!(); + } + } + + let prev_guard = prev_guard.unwrap(); + let prev_node = prev_guard.clone(); + let leaf_lock = { + let pos = last_child_index.unwrap(); + if let Node::Internal(internal) = prev_node { + internal.children[pos].clone() + } else { + unreachable!(); } + }; + + let mut leaf = leaf_lock.write().await; + drop(prev_guard); + let Node::Leaf(leaf_node) = &mut *leaf else { + unreachable!() + }; + + if leaf_node.entries.len() == 2 * self.t - 1 { + return Err(()); } + + match leaf_node.entries.binary_search_by(|(k, _)| k.cmp(&key)) { + Ok(pos) => leaf_node.entries[pos].1 = value, // Обновляем без клонирования + Err(pos) => leaf_node.entries.insert(pos, (key.clone(), value)), + }; + Ok(()) } +} - /// Inserts given value by given key - fn insert(&mut self, key: Rc, value: ChunkHandler, t: usize) -> Option<(Link, Rc)> { - match self { - Node::Leaf(leaf) => { - match leaf.entries.binary_search_by(|(k, _)| k.cmp(&key)) { - Ok(x) => leaf.entries[x] = (key.clone(), value), - Err(x) => leaf.entries.insert(x, (key.clone(), value)), - }; +impl BPlus { + /// Rebuilds links in BPlusTree after loading from file + async fn rebuild_links(&self) { + let leaves = self.collect_leaves().await; + if self.offset.load(Ordering::Acquire) == 0 && self.file_number.load(Ordering::Acquire) == 0 + { + return; + } - if leaf.entries.len() == 2 * t { - return Some(self.split(t)); - } - None - } - Node::Internal(internal_node) => { - let pos = match internal_node.keys.binary_search(&key) { - Ok(x) => x + 1, - Err(x) => x, - }; - let child = internal_node.children[pos].clone(); - let mut borrowed_child = child.borrow_mut(); - let result = borrowed_child.insert(key, value, t); - - match result { - Some((new_child, key)) => { - internal_node.keys.insert(pos, key.clone()); - internal_node.children.insert(pos + 1, new_child); - - match internal_node.keys.len() { - val if val == 2 * t - 1 => Some(self.split(t)), - _ => None, - } + let key_futures: Vec<_> = leaves + .iter() + .map(|leaf| { + let leaf = Arc::clone(leaf); + async move { + let guard = leaf.read().await; + match &*guard { + Node::Leaf(leaf_data) => leaf_data.entries[0].0.clone(), + _ => unreachable!(), } - None => None, } + }) + .collect(); + + let keys = futures::future::join_all(key_futures).await; + + let mut sorted_leaves: Vec<_> = keys.into_iter().zip(leaves.into_iter()).collect(); + + sorted_leaves.sort_by(|(a, _), (b, _)| a.cmp(b)); + + for i in 0..sorted_leaves.len() - 1 { + let current = &sorted_leaves[i].1; + let next = sorted_leaves[i + 1].1.clone(); + + let mut guard = current.write().await; + if let Node::Leaf(leaf) = &mut *guard { + leaf.next = Some(next); } } } - #[allow(unused_variables, dead_code)] - fn remove(&mut self, key: &K, t: usize) -> io::Result<()> { - unimplemented!() - } + /// Collects all leaves from BPlusTree + async fn collect_leaves(&self) -> Vec>>> { + let mut leaves = Vec::new(); + let mut queue = VecDeque::new(); + queue.push_back(self.root.clone()); - /// Gets value from a B+ tree by given key - fn get(&self, key: &K) -> io::Result> { - match self { - Node::Leaf(leaf) => match leaf.entries.binary_search_by(|(k, _)| k.as_ref().cmp(key)) { - Ok(pos) => { - let data_read_result = leaf.entries[pos].1.read()?; - Ok(data_read_result) + while let Some(node) = queue.pop_front() { + let guard = node.read().await; + match &*guard { + Node::Internal(internal) => { + for child in &internal.children { + queue.push_back(child.clone()); + } } - Err(_) => Err(ErrorKind::NotFound.into()), - }, - Node::Internal(internal_node) => { - let pos = match internal_node.keys.binary_search_by(|k| k.as_ref().cmp(key)) { - Ok(x) => x + 1, - Err(x) => x, - }; - - let child = internal_node.children.get(pos); - match child { - Some(x) => x.clone().borrow().get(key), - None => Err(ErrorKind::NotFound.into()), + Node::Leaf(_) => { + leaves.push(node.clone()); } } } + + leaves } -} -impl Database> for BPlus { - fn insert(&mut self, key: K, value: DataContainer<()>) -> io::Result<()> { - match value.extract() { - Data::Chunk(chunk) => self.insert(key, chunk.clone()), - Data::TargetChunk(_chunk) => unimplemented!(), - } + fn open_current_file(path: &Path, number: usize) -> io::Result>> { + Ok(Arc::new(RwLock::new( + File::open(path.join(number.to_string())).unwrap(), + ))) } - fn get(&self, key: &K) -> io::Result> { - self.root.get(key).map(DataContainer::from) + /// Saves this tree by the provided path + pub async fn save(&self, path: &Path) -> io::Result<()> { + let _guard = self.latch.write().await; + let serializable = self.serialize().await; + let file = File::create(path)?; + let writer = BufWriter::new(file); + bincode::serialize_into(writer, &serializable).map_err(io::Error::other) } - fn contains(&self, key: &K) -> bool { - self.root.get(key).is_ok() + /// Loads tree from file by provided path + pub async fn load(path: &Path) -> io::Result { + let file = File::open(path)?; + let reader = BufReader::new(file); + let serializable: SerializableBPlus = + bincode::deserialize_from(reader).map_err(io::Error::other)?; + + Ok(serializable.deserialize().await) } } -/// Iterator for B+ tree -pub struct BPlusIterator { - current_leaf: Option>>>, - current_index: usize, -} +impl Node { + /// Splits node into two and returns new node with it first key + fn split(&mut self, t: usize) -> (Link, Arc) { + match self { + Node::Leaf(leaf) => { + let mut new_leaf_entries = leaf.entries.split_off(t); + new_leaf_entries.reserve_exact(t); + let middle_key = new_leaf_entries[0].0.clone(); -impl Iterator for BPlusIterator { - type Item = (K, Vec); + let new_leaf = Node::Leaf(Leaf { + entries: new_leaf_entries, + next: leaf.next.take(), + }); - fn next(&mut self) -> Option { - loop { - let current = self.current_leaf.as_ref()?.borrow().clone(); - if let Node::Leaf(leaf) = current { - if self.current_index < leaf.entries.len() { - let (key, handler) = &leaf.entries[self.current_index]; - let value = handler.read().ok()?; - let key = Rc::unwrap_or_clone(key.clone()); - self.current_index += 1; - return Some((key, value)); - } else { - match &leaf.next { - Some(next_leaf) => { - self.current_leaf = Some(Rc::clone(next_leaf)); - self.current_index = 0; - } - None => { - self.current_leaf = None; - return None; - } - } - } - } else { - return None; + let new_leaf_link = Arc::new(RwLock::new(new_leaf)); + leaf.next = Some(new_leaf_link.clone()); + + (new_leaf_link, middle_key) } - } - } -} + Node::Internal(internal_node) => { + let mut new_node_keys = internal_node.keys.split_off(t - 1); + let middle_key = new_node_keys.remove(0); -impl IntoIterator for BPlus { - type Item = (K, Vec); - type IntoIter = BPlusIterator; - fn into_iter(self) -> Self::IntoIter { - let mut current = Some(Rc::new(RefCell::new(self.root))); - while let Some(node) = current.clone() { - let borrowed = node.borrow(); - match &*borrowed { - Node::Internal(internal) => { - current = Some(Rc::clone(&internal.children[0])); - } - Node::Leaf(_) => break, + let mut new_node_children = internal_node.children.split_off(t); + new_node_keys.reserve_exact(t); + new_node_children.reserve_exact(t); + + let new_node = Node::Internal(InternalNode { + children: new_node_children, + keys: new_node_keys, + }); + + (Arc::new(RwLock::new(new_node)), middle_key) } } + } - BPlusIterator { - current_leaf: current, - current_index: 0, - } + #[allow(unused_variables, dead_code)] + fn remove(&mut self, key: &K, t: usize) -> io::Result<()> { + unimplemented!() } } #[cfg(test)] mod tests { - use tempdir::TempDir; - use super::*; + use tempfile::TempDir; - #[test] - fn test_node_split_leaf() { - let tempdir = TempDir::new("split_leaf").unwrap(); - let mut tree: BPlus = BPlus::new(2, tempdir.path().into()).unwrap(); - - for i in 1..5 { - tree.insert(i, vec![i as u8]).unwrap(); - } - - if let Node::Internal(root) = &tree.root { - assert_eq!(root.keys.len(), 1); - assert_eq!(root.children.len(), 2); - } else { - assert!(false); - } + fn create_test_tree(t: usize, name: &str) -> (BPlus, TempDir) { + let temp_dir = TempDir::with_prefix(name).unwrap(); + let tree = BPlus::new(t, temp_dir.path().to_path_buf()).unwrap(); + (tree, temp_dir) } - #[test] - fn test_file_rotation() { - let tempdir = TempDir::new("file_rotation").unwrap(); - let mut tree = BPlus::new(2, tempdir.path().into()).unwrap(); - tree.max_file_size = 128; + #[tokio::test(flavor = "multi_thread")] + async fn test_multiple_inserts() { + let (tree, _temp) = create_test_tree(2, "multiple_inserts"); - let data = vec![0; 200]; - tree.insert(1, data).unwrap(); - - let small_data = vec![0; 10]; - tree.insert(2, small_data).unwrap(); + for i in 1..=4 { + tree.insert(i, vec![i as u8]).await; + } - assert_eq!(tree.file_number, 1); - assert_eq!(tree.offset, 10); + for i in 1..=4 { + let result = tree.get(&i).await.unwrap(); + assert_eq!(result, vec![i as u8]); + } } - #[test] - fn test_leaf_linking() { - let tempdir = TempDir::new("leaf_link").unwrap(); - let mut tree: BPlus = BPlus::new(2, tempdir.path().into()).unwrap(); + #[tokio::test(flavor = "multi_thread")] + async fn test_concurrent_inserts() { + let (tree, _temp) = create_test_tree(2, "concurrent_inserts"); + let tree = Arc::new(tokio::sync::RwLock::new(tree)); + + let mut handles = vec![]; + for i in 0..50 { + let tree = tree.clone(); + handles.push(tokio::spawn(async move { + let tree = tree.write().await; + tree.insert(i, vec![i as u8]).await; + })); + } - for i in 1..=9 { - tree.insert(i, vec![i as u8]).unwrap(); + for handle in handles { + handle.await.unwrap(); } - let mut current = Some(Rc::new(RefCell::new(tree.root))); - while let Some(node) = current.clone() { - let borrowed = node.borrow(); - match &*borrowed { - Node::Internal(internal) => { - current = Some(Rc::clone(&internal.children[0])); - } - Node::Leaf(_) => break, - } + let tree = tree.read().await; + for i in 0..50 { + let result = tree.get(&i).await.unwrap(); + assert_eq!(result, vec![i as u8]); } + } - let mut collected = vec![]; - if let Node::Leaf(first_leaf) = &*current.clone().as_ref().unwrap().borrow() { - let mut current_leaf = Some(Rc::new(RefCell::new(Node::Leaf(first_leaf.clone())))); + #[tokio::test(flavor = "multi_thread")] + async fn test_root_split() { + let (tree, _temp) = create_test_tree(2, "root_split"); - while let Some(leaf_ref) = current_leaf { - let leaf_guard = leaf_ref.borrow(); + tree.insert(1, vec![1]).await; + tree.insert(2, vec![2]).await; + tree.insert(3, vec![3]).await; + tree.insert(4, vec![4]).await; - if let Node::Leaf(leaf) = &*leaf_guard { - collected.extend(leaf.entries.iter().map(|(k, _)| **k)); - current_leaf = leaf.next.as_ref().map(Rc::clone); - } else { - break; - } + let root = tree.root.read().await; + match &*root { + Node::Internal(internal) => { + assert_eq!(internal.keys.len(), 1); + assert_eq!(internal.children.len(), 2); } + _ => panic!("Root should be internal node after split"), } - - assert_eq!(collected, (1..=9).collect::>()); - assert_eq!(collected.len(), 9); } - #[test] - fn test_consecutive_file_handling() { - let tempdir = TempDir::new("consecutive").unwrap(); - let mut tree = BPlus::new(2, tempdir.path().into()).unwrap(); - tree.max_file_size = 256; + #[tokio::test(flavor = "multi_thread")] + async fn test_large_value_storage() { + let temp_dir = TempDir::new().unwrap(); + let mut tree = BPlus::new(2, temp_dir.path().to_path_buf()).unwrap(); + tree.max_file_size = 100; - for i in 0..500 { - tree.insert(i, vec![i as u8; 100]).unwrap(); - } + let large_data = vec![7; 150]; + tree.insert(1, large_data.clone()).await; - assert!(tree.file_number > 1); + let result = tree.get(&1).await.unwrap(); + assert_eq!(result, large_data); + tree.insert(2, large_data.clone()).await; + let result = tree.get(&1).await.unwrap(); + assert_eq!(result, large_data); + + assert!( + tree.file_number.load(std::sync::atomic::Ordering::SeqCst) >= 1, + "Should create multiple files" + ); } - #[test] - fn test_node_consistency_after_splits() { - let tempdir = TempDir::new("consistency").unwrap(); - let mut tree = BPlus::new(2, tempdir.path().into()).unwrap(); + #[tokio::test] + async fn test_save_load_empty_tree() { + let tempdir = TempDir::new().unwrap(); + let tree_path = tempdir.path().join("empty_tree.bin"); - for i in 0..10 { - tree.insert(i, vec![i as u8]).unwrap(); - } + let tree = BPlus::::new(2, tempdir.path().into()).unwrap(); - if let Node::Internal(root) = &tree.root { - assert!(!root.keys.is_empty()); - assert!(root.children.len() >= 2); + tree.save(&tree_path).await.unwrap(); - for key in &root.keys { - let key_val = **key; - assert!(key_val > 0 && key_val < 10); - } - } + let loaded_tree = BPlus::::load(&tree_path).await.unwrap(); + + assert_eq!(tree.t, loaded_tree.t); + assert_eq!(tree.path, loaded_tree.path); + assert_eq!( + tree.file_number.load(Ordering::SeqCst), + loaded_tree.file_number.load(Ordering::SeqCst) + ); + assert_eq!( + tree.offset.load(Ordering::SeqCst), + loaded_tree.offset.load(Ordering::SeqCst) + ); + assert!(loaded_tree.get(&42).await.is_err()); } } diff --git a/tests/bplus_tests.rs b/tests/bplus_tests.rs index 1a6609e..cfa2a99 100644 --- a/tests/bplus_tests.rs +++ b/tests/bplus_tests.rs @@ -1,320 +1,330 @@ extern crate chunkfs; +use bplus_tree::bplus_tree::BPlus; use std::collections::HashMap; use std::path::PathBuf; - -use bplus_tree::bplus_tree::BPlus; -use rand::seq::SliceRandom; use tempdir::TempDir; -#[test] -fn test_non_existent_key() { +#[tokio::test(flavor = "multi_thread")] +async fn test_non_existent_key() { let tempdir = TempDir::new("non_existent").unwrap(); - let mut tree: BPlus = BPlus::new(2, tempdir.path().into()).unwrap(); - tree.insert(1, vec![1]).unwrap(); - assert!(tree.get(&2).is_err()); + let tree: BPlus = BPlus::new(2, tempdir.path().into()).unwrap(); + tree.insert(1, vec![1]).await; + assert!(tree.get(&2).await.is_err()); } -#[test] -fn test_overwrite_existing_key() { +#[tokio::test(flavor = "multi_thread")] +async fn test_overwrite_existing_key() { let tempdir = TempDir::new("overwrite").unwrap(); - let mut tree: BPlus = BPlus::new(2, tempdir.path().into()).unwrap(); + let tree: BPlus = BPlus::new(2, tempdir.path().into()).unwrap(); - tree.insert(1, vec![1]).unwrap(); - tree.insert(1, vec![42]).unwrap(); + tree.insert(1, vec![1]).await; + tree.insert(1, vec![42]).await; - assert_eq!(tree.get(&1).unwrap(), vec![42]); + assert_eq!(tree.get(&1).await.unwrap(), vec![42]); } -#[test] -fn test_insert_and_find() { +#[tokio::test(flavor = "multi_thread")] +async fn test_insert_and_find() { let tempdir = TempDir::new("1").unwrap(); let path = PathBuf::new().join(tempdir.path()); - let mut tree: BPlus = BPlus::new(2, path).unwrap(); + let tree: BPlus = BPlus::new(2, path).unwrap(); for i in 1..6 { - let _ = tree.insert(i, vec![i as u8; 1]); + tree.insert(i, vec![i as u8; 1]).await; } for i in 1..6 { - let a = tree.get(&i).unwrap(); + let a = tree.get(&i).await.unwrap(); assert_eq!(a, vec![i as u8; 1]); } } -#[test] -fn test_insert_and_find_many_nodes() { +#[tokio::test(flavor = "multi_thread")] +async fn test_insert_and_find_many_nodes() { let tempdir = TempDir::new("4").unwrap(); let path = PathBuf::new().join(tempdir.path()); - let mut tree: BPlus = BPlus::new(2, path).unwrap(); + let tree: BPlus = BPlus::new(2, path).unwrap(); for i in 1..255 { - let _ = tree.insert(i, vec![i as u8; 1]); + tree.insert(i, vec![i as u8; 1]).await; } for i in 1..255 { - assert_eq!(tree.get(&(i as usize)).unwrap(), vec![i as u8; 1]); + assert_eq!(tree.get(&(i as usize)).await.unwrap(), vec![i as u8; 1]); } } -#[test] -fn test_large_data_consecutive_numbers() { +#[tokio::test(flavor = "multi_thread")] +async fn test_large_data_consecutive_numbers() { let tempdir = TempDir::new("6").unwrap(); let path = PathBuf::new().join(tempdir.path()); - let mut tree: BPlus = BPlus::new(100, path).unwrap(); + let tree: BPlus = BPlus::new(100, path).unwrap(); for i in 1..10000 { - let _ = tree.insert(i, vec![i as u8; 1064]); + tree.insert(i, vec![i as u8; 1064]).await; } for i in 1..10000 { - let a = tree.get(&i).unwrap(); + let a = tree.get(&i).await.unwrap(); assert_eq!(a, vec![i as u8; 1064]); } } -#[test] -fn test_large_data() { +#[tokio::test(flavor = "multi_thread")] +async fn test_large_data() { let tempdir = TempDir::new("7").unwrap(); let path = PathBuf::new().join(tempdir.path()); - let mut tree: BPlus = BPlus::new(2, path).unwrap(); + let tree: BPlus = BPlus::new(2, path).unwrap(); let mut htable = HashMap::>::new(); for i in 1..10000 { - let key; - key = i * 113; - let _ = tree.insert(key, vec![key as u8; 1064]); + let key = i * 113; + tree.insert(key, vec![key as u8; 1064]).await; htable.insert(key, vec![key as u8; 1064]); } for (key, value) in htable { - assert_eq!(tree.get(&key).unwrap(), value); + assert_eq!(tree.get(&key).await.unwrap(), value); } } -#[test] -fn test_couple_of_same_keys_inserted() { +#[tokio::test(flavor = "multi_thread")] +async fn test_couple_of_same_keys_inserted() { let tempdir = TempDir::new("8").unwrap(); - let mut tree: BPlus = BPlus::new(2, PathBuf::new().join(tempdir.path())).unwrap(); + let tree: BPlus = BPlus::new(2, PathBuf::new().join(tempdir.path())).unwrap(); for i in 1..100 { - tree.insert(i, vec![1u8]).unwrap(); + tree.insert(i, vec![1u8]).await; } for i in 1..100 { for j in 1..100 { - tree.insert(i, vec![j as u8]).unwrap(); + tree.insert(i, vec![j as u8]).await; } } for i in 1..100 { for _ in 1..100 { - tree.get(&i).unwrap(); + tree.get(&i).await.unwrap(); } } } -#[test] -fn test_same_keys_inserted() { +#[tokio::test(flavor = "multi_thread")] +async fn test_same_keys_inserted() { let tempdir = TempDir::new("10").unwrap(); - let mut tree: BPlus = BPlus::new(2, PathBuf::new().join(tempdir.path())).unwrap(); + let tree: BPlus = BPlus::new(2, PathBuf::new().join(tempdir.path())).unwrap(); let mut keys = vec![]; - for _ in 1..1000 { + for _ in 1..100000 { let key: usize = rand::random::() % 10000; keys.push(key); } for key in keys.clone() { - tree.insert(key, vec![key as u8]).unwrap(); + tree.insert(key, vec![key as u8]).await; } for key in keys { - assert_eq!(vec![key as u8], tree.get(&key).unwrap()); + assert_eq!(vec![key as u8], tree.get(&key).await.unwrap()); } let key: usize = rand::random(); - tree.insert(key, vec![0u8]).unwrap(); + tree.insert(key, vec![0u8]).await; for i in 1..255 { - assert_eq!(vec![i - 1u8], tree.get(&key).unwrap()); - tree.insert(key, vec![i]).unwrap(); + assert_eq!(vec![i - 1u8], tree.get(&key).await.unwrap()); + tree.insert(key, vec![i]).await; } } -#[test] -fn test_iterator() { - let tempdir = TempDir::new("iterator_test").unwrap(); - let mut tree: BPlus = BPlus::new(2, tempdir.path().into()).unwrap(); +#[tokio::test(flavor = "multi_thread")] +async fn test_same_10k_keys_inserted() { + let tempdir = TempDir::new("same_10k_keys").unwrap(); + let tree: BPlus = BPlus::new(2, PathBuf::new().join(tempdir.path())).unwrap(); - for i in 1..5 { - tree.insert(i, vec![i as u8]).unwrap(); + for i in 0..10000 { + tree.insert(i, vec![i as u8]).await; } - let mut iter = tree.into_iter(); - assert_eq!(iter.next(), Some((1, vec![1]))); - assert_eq!(iter.next(), Some((2, vec![2]))); - assert_eq!(iter.next(), Some((3, vec![3]))); - assert_eq!(iter.next(), Some((4, vec![4]))); - assert_eq!(iter.next(), None); + for i in 0..10000 { + tree.insert(i, vec![i as u8]).await; + } + + for key in 1..10000 { + assert_eq!(vec![key as u8], tree.get(&key).await.unwrap()); + } } -#[test] -fn test_empty_tree() { + +#[tokio::test(flavor = "multi_thread")] +async fn test_empty_tree() { let tempdir = TempDir::new("empty").unwrap(); let tree: BPlus = BPlus::new(2, tempdir.path().into()).unwrap(); - assert!(tree.get(&1).is_err()); + assert!(tree.get(&1).await.is_err()); } -#[test] -fn test_single_entry() { +#[tokio::test(flavor = "multi_thread")] +async fn test_single_entry() { let tempdir = TempDir::new("single").unwrap(); - let mut tree = BPlus::new(2, tempdir.path().into()).unwrap(); - tree.insert(42, vec![1, 2, 3]).unwrap(); - assert_eq!(tree.get(&42).unwrap(), vec![1, 2, 3]); + let tree = BPlus::new(2, tempdir.path().into()).unwrap(); + tree.insert(42, vec![1, 2, 3]).await; + assert_eq!(tree.get(&42).await.unwrap(), vec![1, 2, 3]); } -#[test] -fn test_reverse_order_insert() { +#[tokio::test(flavor = "multi_thread")] +async fn test_reverse_order_insert() { let tempdir = TempDir::new("reverse").unwrap(); - let mut tree: BPlus = BPlus::new(3, tempdir.path().into()).unwrap(); + let tree: BPlus = BPlus::new(3, tempdir.path().into()).unwrap(); for i in (1..100).rev() { - tree.insert(i, vec![i as u8]).unwrap(); + tree.insert(i, vec![i as u8]).await; } for i in 1..100 { - assert_eq!(tree.get(&i).unwrap(), vec![i as u8]); - } -} - -#[test] -fn test_iterator_empty() { - let tempdir = TempDir::new("iter_empty").unwrap(); - let tree: BPlus = BPlus::new(2, tempdir.path().into()).unwrap(); - let mut iter = tree.into_iter(); - assert_eq!(iter.next(), None); -} - -#[test] -fn test_iterator_single() { - let tempdir = TempDir::new("iter_single").unwrap(); - let mut tree = BPlus::new(2, tempdir.path().into()).unwrap(); - tree.insert(1, vec![1]).unwrap(); - let mut iter = tree.into_iter(); - assert_eq!(iter.next(), Some((1, vec![1]))); - assert_eq!(iter.next(), None); -} - -#[test] -fn test_iterator_large_dataset() { - let tempdir = TempDir::new("iter_large").unwrap(); - let mut tree = BPlus::new(10, tempdir.path().into()).unwrap(); - let mut expected = Vec::new(); - - for i in 1..=21 { - tree.insert(i, vec![i as u8]).unwrap(); - expected.push((i, vec![i as u8])); + assert_eq!(tree.get(&i).await.unwrap(), vec![i as u8]); } - - let result: Vec<_> = tree.into_iter().collect(); - assert_eq!(result, expected); } -#[test] -fn test_minimal_degree() { +#[tokio::test(flavor = "multi_thread")] +async fn test_minimal_degree() { let tempdir = TempDir::new("min_degree").unwrap(); - let mut tree = BPlus::new(1, tempdir.path().into()).unwrap(); + let tree = BPlus::new(1, tempdir.path().into()).unwrap(); for i in 1..=10 { - tree.insert(i, vec![i as u8]).unwrap(); + tree.insert(i, vec![i as u8]).await; } - assert_eq!(tree.get(&5).unwrap(), vec![5]); + assert_eq!(tree.get(&5).await.unwrap(), vec![5]); } -#[test] -fn test_key_duplication() { +#[tokio::test(flavor = "multi_thread")] +async fn test_key_duplication() { let tempdir = TempDir::new("dupes").unwrap(); - let mut tree = BPlus::new(2, tempdir.path().into()).unwrap(); + let tree = BPlus::new(2, tempdir.path().into()).unwrap(); for _ in 0..10 { - tree.insert(42, vec![1]).unwrap(); - tree.insert(42, vec![2]).unwrap(); + tree.insert(42, vec![1]).await; + tree.insert(42, vec![2]).await; } - assert_eq!(tree.get(&42).unwrap(), vec![2]); + assert_eq!(tree.get(&42).await.unwrap(), vec![2]); } -#[test] -fn test_string_keys() { +#[tokio::test(flavor = "multi_thread")] +async fn test_string_keys() { let tempdir = TempDir::new("string_keys").unwrap(); - let mut tree = BPlus::new(2, tempdir.path().into()).unwrap(); + let tree = BPlus::new(2, tempdir.path().into()).unwrap(); - tree.insert("apple".to_string(), b"fruit".to_vec()).unwrap(); - tree.insert("banana".to_string(), b"yellow".to_vec()) - .unwrap(); + tree.insert("apple".to_string(), b"fruit".to_vec()).await; + tree.insert("banana".to_string(), b"yellow".to_vec()).await; - assert_eq!(tree.get(&"apple".to_string()).unwrap(), b"fruit"); - assert_eq!(tree.get(&"banana".to_string()).unwrap(), b"yellow"); + assert_eq!(tree.get(&"apple".to_string()).await.unwrap(), b"fruit"); + assert_eq!(tree.get(&"banana".to_string()).await.unwrap(), b"yellow"); } -#[test] -fn test_stress_1m_entries() { +#[tokio::test(flavor = "multi_thread")] +async fn test_stress_1m_entries() { let tempdir = TempDir::new("stress_1m").unwrap(); - let mut tree = BPlus::new(100, tempdir.path().into()).unwrap(); + let tree = BPlus::new(100, tempdir.path().into()).unwrap(); for i in 0..1_000_000 { - tree.insert(i, vec![i as u8]).unwrap(); + tree.insert(i, vec![i as u8]).await; } for i in 0..1_000_000 { - assert_eq!(tree.get(&i).unwrap(), vec![i as u8]); + assert_eq!(tree.get(&i).await.unwrap(), vec![i as u8]); + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 16)] +async fn test_stress_1m_entries_concurrent() { + use std::sync::Arc; + use tokio::task; + + let tempdir = TempDir::new("stress_1m").unwrap(); + let tree = Arc::new(BPlus::new(100, tempdir.path().into()).unwrap()); + let num_tasks = 1000; + let entries_per_task = 1000; + let barrier = Arc::new(tokio::sync::Barrier::new(num_tasks + 1)); + + let mut insert_handles = Vec::with_capacity(num_tasks); + for task_id in 0..num_tasks { + let tree = Arc::clone(&tree); + let barrier = Arc::clone(&barrier); + + insert_handles.push(task::spawn(async move { + barrier.wait().await; + + for i in 0..entries_per_task { + let key = (task_id * entries_per_task) + i; + tree.insert(key, vec![key as u8]).await; + } + })); + } + + barrier.wait().await; + + for handle in insert_handles { + handle.await.unwrap(); + } + + let mut verify_handles = Vec::with_capacity(num_tasks); + for task_id in 0..num_tasks { + let tree = Arc::clone(&tree); + + verify_handles.push(task::spawn(async move { + for i in 0..entries_per_task { + let key = (task_id * entries_per_task) + i; + let value = tree.get(&key).await.unwrap(); + assert_eq!(value, vec![key as u8], "Invalid value for key {}", key); + } + })); + } + + for handle in verify_handles { + handle.await.unwrap(); } } -#[test] -fn test_find_nonexistent_after_splits() { +#[tokio::test(flavor = "multi_thread")] +async fn test_find_nonexistent_after_splits() { let tempdir = TempDir::new("nonexistent").unwrap(); - let mut tree = BPlus::new(2, tempdir.path().into()).unwrap(); + let tree = BPlus::new(2, tempdir.path().into()).unwrap(); for i in 0..1000 { - tree.insert(i, vec![1]).unwrap(); + tree.insert(i, vec![1]).await; } - assert!(tree.get(&1001).is_err()); + assert!(tree.get(&1001).await.is_err()); } +#[tokio::test] +async fn test_save_load_small_tree() { + let tempdir = TempDir::new("saveload_small").unwrap(); + let tree_path = tempdir.path().join("small_tree.bin"); -#[test] -fn test_entry_ordering() { - let tempdir = TempDir::new("ordering").unwrap(); - let mut tree = BPlus::new(3, tempdir.path().into()).unwrap(); - let mut rng = rand::thread_rng(); - let mut nums: Vec = (0..1000).collect(); - nums.shuffle(&mut rng); + let tree = BPlus::::new(2, tempdir.path().into()).unwrap(); + tree.insert(10, vec![1, 2, 3]).await; + tree.insert(20, vec![4, 5, 6]).await; + tree.insert(5, vec![0]).await; - for &num in &nums { - tree.insert(num, vec![num as u8]).unwrap(); - } + tree.save(&tree_path).await.unwrap(); - let mut sorted = nums.clone(); - sorted.sort_unstable(); - let result: Vec<_> = tree.into_iter().map(|(k, _)| k).collect(); - assert_eq!(result, sorted); + let loaded_tree = BPlus::::load(&tree_path).await.unwrap(); + + assert_eq!(loaded_tree.get(&10).await.unwrap(), vec![1, 2, 3]); + assert_eq!(loaded_tree.get(&20).await.unwrap(), vec![4, 5, 6]); + assert_eq!(loaded_tree.get(&5).await.unwrap(), vec![0]); + assert!(loaded_tree.get(&99).await.is_err()); } -#[test] -fn test_custom_data_types() { - #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Default)] - struct ComplexKey { - id: u64, - name: String, - } +#[tokio::test] +async fn test_save_load_large_tree() { + let tempdir = TempDir::new("large_load_save").unwrap(); + let tree_path = tempdir.path().join("large_tree.bin"); + let mut tree = BPlus::::new(2, tempdir.path().into()).unwrap(); - let tempdir = TempDir::new("custom_type").unwrap(); - let mut tree = BPlus::new(2, tempdir.path().into()).unwrap(); + for i in 0..100000 { + tree.insert(i, vec![(i % 256) as u8; 200]).await; + } + tree.save(&tree_path).await.unwrap(); - let key1 = ComplexKey { - id: 1, - name: "A".to_string(), - }; - let key2 = ComplexKey { - id: 2, - name: "B".to_string(), - }; + let loaded_tree = BPlus::::load(&tree_path).await.unwrap(); - tree.insert(key1.clone(), b"data1".to_vec()).unwrap(); - tree.insert(key2.clone(), b"data2".to_vec()).unwrap(); + for key in 0..100000 { + let expected = vec![(key % 256) as u8; 200]; + assert_eq!(loaded_tree.get(&key).await.unwrap(), expected); + } - assert_eq!(tree.get(&key1).unwrap(), b"data1"); - assert_eq!(tree.get(&key2).unwrap(), b"data2"); + assert!(loaded_tree.get(&100_000).await.is_err()); } diff --git a/tests/chunkfs_tests.rs b/tests/chunkfs_tests.rs index 23c4b89..2795f7c 100644 --- a/tests/chunkfs_tests.rs +++ b/tests/chunkfs_tests.rs @@ -7,19 +7,23 @@ use std::path::PathBuf; use approx::assert_relative_eq; -use bplus_tree::bplus_tree::BPlus; +use bplus_tree::bplus_tree::BPlusStorage; use chunkfs::chunkers::{FSChunker, LeapChunker}; use chunkfs::hashers::SimpleHasher; use chunkfs::{create_cdc_filesystem, DataContainer, Database, WriteMeasurements}; use tempdir::TempDir; +use tokio::runtime::Builder; + const MB: usize = 1024 * 1024; #[test] fn write_read_complete_test() { let tempdir = &TempDir::new("storage1").unwrap(); let path = PathBuf::new().join(tempdir.path()); - let mut fs = create_cdc_filesystem(BPlus::new(100, path).unwrap(), SimpleHasher); + let runtime = Builder::new_multi_thread().enable_all().build().unwrap(); + let mut fs = + create_cdc_filesystem(BPlusStorage::new(runtime, 100, path).unwrap(), SimpleHasher); let mut handle = fs.create_file("file", LeapChunker::default()).unwrap(); fs.write_to_file(&mut handle, &[1; MB]).unwrap(); @@ -38,8 +42,9 @@ fn write_read_complete_test() { fn write_read_blocks_test() { let tempdir = &TempDir::new("storage2").unwrap(); let path = PathBuf::new().join(tempdir.path()); - let mut fs = create_cdc_filesystem(BPlus::new(100, path).unwrap(), SimpleHasher); - + let runtime = Builder::new_multi_thread().enable_all().build().unwrap(); + let mut fs = + create_cdc_filesystem(BPlusStorage::new(runtime, 100, path).unwrap(), SimpleHasher); let mut handle = fs.create_file("file", FSChunker::new(4096)).unwrap(); let ones = vec![1; MB]; @@ -61,8 +66,9 @@ fn write_read_blocks_test() { fn read_file_with_size_less_than_1mb() { let tempdir = &TempDir::new("storage3").unwrap(); let path = PathBuf::new().join(tempdir.path()); - let mut fs = create_cdc_filesystem(BPlus::new(100, path).unwrap(), SimpleHasher); - + let runtime = Builder::new_multi_thread().enable_all().build().unwrap(); + let mut fs = + create_cdc_filesystem(BPlusStorage::new(runtime, 100, path).unwrap(), SimpleHasher); let mut handle = fs.create_file("file", FSChunker::new(4096)).unwrap(); let ones = vec![1; 10]; @@ -78,7 +84,9 @@ fn read_file_with_size_less_than_1mb() { fn write_read_big_file_at_once() { let tempdir = &TempDir::new("storage4").unwrap(); let path = PathBuf::new().join(tempdir.path()); - let mut fs = create_cdc_filesystem(BPlus::new(100, path).unwrap(), SimpleHasher); + let runtime = Builder::new_multi_thread().enable_all().build().unwrap(); + let mut fs = + create_cdc_filesystem(BPlusStorage::new(runtime, 100, path).unwrap(), SimpleHasher); let mut handle = fs.create_file("file", FSChunker::new(4096)).unwrap(); @@ -94,7 +102,9 @@ fn write_read_big_file_at_once() { fn two_file_handles_to_one_file() { let tempdir = &TempDir::new("storage6").unwrap(); let path = PathBuf::new().join(tempdir.path()); - let mut fs = create_cdc_filesystem(BPlus::new(100, path).unwrap(), SimpleHasher); + let runtime = Builder::new_multi_thread().enable_all().build().unwrap(); + let mut fs = + create_cdc_filesystem(BPlusStorage::new(runtime, 100, path).unwrap(), SimpleHasher); let mut handle1 = fs.create_file("file", LeapChunker::default()).unwrap(); let mut handle2 = fs.open_file("file", LeapChunker::default()).unwrap(); fs.write_to_file(&mut handle1, &[1; MB]).unwrap(); @@ -160,7 +170,9 @@ fn dedup_ratio_is_correct_for_fixed_size_chunker() { fn readonly_file_handle_cannot_write_can_read() { let tempdir = &TempDir::new("storage8").unwrap(); let path = PathBuf::new().join(tempdir.path()); - let mut fs = create_cdc_filesystem(BPlus::new(100, path).unwrap(), SimpleHasher); + let runtime = Builder::new_multi_thread().enable_all().build().unwrap(); + let mut fs = + create_cdc_filesystem(BPlusStorage::new(runtime, 100, path).unwrap(), SimpleHasher); let mut fh = fs.create_file("file", FSChunker::default()).unwrap(); fs.write_to_file(&mut fh, &[1; MB]).unwrap(); fs.close_file(fh).unwrap(); @@ -189,7 +201,9 @@ fn readonly_file_handle_cannot_write_can_read() { fn write_from_stream_slice() { let tempdir = &TempDir::new("storage9").unwrap(); let path = PathBuf::new().join(tempdir.path()); - let mut fs = create_cdc_filesystem(BPlus::new(100, path).unwrap(), SimpleHasher); + let runtime = Builder::new_multi_thread().enable_all().build().unwrap(); + let mut fs = + create_cdc_filesystem(BPlusStorage::new(runtime, 100, path).unwrap(), SimpleHasher); let mut fh = fs.create_file("file", FSChunker::default()).unwrap(); fs.write_from_stream(&mut fh, &[1; MB * 2][..]).unwrap(); fs.close_file(fh).unwrap(); @@ -208,7 +222,9 @@ fn write_from_stream_buf_reader() { let tempdir = &TempDir::new("storage10").unwrap(); let path = PathBuf::new().join(tempdir.path()); - let mut fs = create_cdc_filesystem(BPlus::new(100, path).unwrap(), SimpleHasher); + let runtime = Builder::new_multi_thread().enable_all().build().unwrap(); + let mut fs = + create_cdc_filesystem(BPlusStorage::new(runtime, 100, path).unwrap(), SimpleHasher); let mut fh = fs.create_file("file", FSChunker::default()).unwrap(); fs.write_from_stream(&mut fh, file).unwrap();