From 655ec8f62edbe872f8ae6dd4787ec9dd97ecd5f5 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Sun, 4 May 2025 21:44:49 +0300 Subject: [PATCH 01/22] Rewrite get and inset from recursive to iterative. Change root type from Node to Link --- src/bplus_tree.rs | 244 ++++++++++++++++++++++++++++------------------ 1 file changed, 151 insertions(+), 93 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index d5c1652..5ac18eb 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -75,7 +75,7 @@ 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. @@ -97,27 +97,35 @@ impl fmt::Display for BPlus { } impl BPlus { - fn fmt_node(node: &Node, level: usize, f: &mut fmt::Formatter<'_>) -> fmt::Result { - match node { + fn fmt_node(node: &Link, level: usize, f: &mut fmt::Formatter<'_>) -> fmt::Result { + let node_ref = node.borrow(); + + match &*node_ref { Node::Internal(internal) => { writeln!( f, "{}[Internal] keys: {:?}", " ".repeat(level), - internal.keys + internal.keys.iter().map(|k| k.as_ref()).collect::>() )?; for child in &internal.children { - BPlus::fmt_node(&child.borrow(), level + 1, f)?; + Self::fmt_node(child, level + 1, f)?; } Ok(()) } Node::Leaf(leaf) => { writeln!( f, - "{}[Leaf] entries: {:?}", + "{}[Leaf] entries: {:?}, next: {}", " ".repeat(level), leaf.entries + .iter() + .map(|(k, v)| (k.as_ref(), v)) + .collect::>(), + leaf.next + .as_ref() + .map_or("None".into(), |n| format!("{:p}", Rc::as_ptr(n))) ) } } @@ -134,7 +142,7 @@ impl BPlus { create_dir_all(&path)?; let current_file = File::create(path_to_file)?; Ok(Self { - root: Node::Leaf(Leaf::default()), + root: Rc::new(RefCell::new(Node::Leaf(Leaf::default()))), t, path, file_number: 0, @@ -171,15 +179,16 @@ impl BPlus { let key = Rc::new(key); - 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 { + if let Some((new_node_link, new_key)) = + Node::insert(self.root.clone(), key, value_to_insert, self.t) + { + //let mut prev_root = self.root.clone(); + if let Node::Leaf(ref mut leaf) = &mut *(self.root.clone()).borrow_mut() { leaf.next = Some(new_node_link.clone()); - prev_root = Node::Leaf(leaf); } 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(self.root.clone()); new_root_children.push(new_node_link); new_root_keys.push(new_key); let new_root = Node::::Internal(InternalNode { @@ -187,8 +196,9 @@ impl BPlus { keys: new_root_keys, }); - self.root = new_root; + self.root = Rc::new(RefCell::new(new_root)); } + Ok(()) } @@ -199,7 +209,7 @@ impl BPlus { /// Gets value from a B+ tree by given key pub fn get(&self, key: &K) -> io::Result> { - self.root.get(key) + Node::get(self.root.clone(), key) } } @@ -240,43 +250,70 @@ impl Node { } } - /// 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)), - }; + /// Inserts given value by given key + fn insert( + root: Link, + key: Rc, + value: ChunkHandler, + t: usize, + ) -> Option<(Link, Rc)> { + let mut path = Vec::new(); // Path to leaf + let mut current = root; + let mut split_result; + + // Descent to the leaf + loop { + let current_clone = current.clone(); + let mut current_node = current_clone.borrow_mut(); + 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 * t { + Some(current_node.split(t)) + } else { + None + }; + break; + } + Node::Internal(internal) => { + let pos = match internal.keys.binary_search(&key) { + Ok(pos) => pos + 1, + Err(pos) => pos, + }; + + let next_node = internal.children[pos].clone(); - if leaf.entries.len() == 2 * t { - return Some(self.split(t)); + path.push((current, pos)); + + current = next_node; } - 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, - } + } + + // Going up to the root splitting nodes if needed + while let Some((parent, pos)) = path.pop() { + if let Some((new_node, median)) = split_result.take() { + let mut node = parent.borrow_mut(); + 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 * t - 1 { + split_result = Some(node.split(t)); + } else { + split_result = None; } - None => None, + } else { + break; // No need to split the nodes anymore } } } + + split_result } #[allow(unused_variables, dead_code)] @@ -285,27 +322,37 @@ impl Node { } /// 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) + fn get(root: Link, key: &K) -> io::Result> { + let mut current = root; + + loop { + let current_clone = current.clone(); + let node = current_clone.borrow(); + + 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()?; + Ok(data_read_result) + } + Err(_) => Err(ErrorKind::NotFound.into()), + }; } - 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::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 => return Err(ErrorKind::NotFound.into()), + }; } } + + drop(node); } } } @@ -319,11 +366,11 @@ impl Database> for BPlus< } fn get(&self, key: &K) -> io::Result> { - self.root.get(key).map(DataContainer::from) + Node::get(self.root.clone(), key).map(DataContainer::from) } fn contains(&self, key: &K) -> bool { - self.root.get(key).is_ok() + Node::get(self.root.clone(), key).is_ok() } } @@ -369,7 +416,7 @@ 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))); + let mut current = Some(self.root); while let Some(node) = current.clone() { let borrowed = node.borrow(); match &*borrowed { @@ -396,17 +443,28 @@ mod tests { #[test] fn test_node_split_leaf() { let tempdir = TempDir::new("split_leaf").unwrap(); - let mut tree: BPlus = BPlus::new(2, tempdir.path().into()).unwrap(); + let mut tree = 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); + let root = tree.root.borrow(); + if let Node::Internal(root_node) = &*root { + assert_eq!(root_node.keys.len(), 1); + assert_eq!(root_node.children.len(), 2); + + let left_child = root_node.children[0].borrow(); + if let Node::Leaf(left_leaf) = &*left_child { + assert_eq!(left_leaf.entries.len(), 2); + } + + let right_child = root_node.children[1].borrow(); + if let Node::Leaf(right_leaf) = &*right_child { + assert_eq!(right_leaf.entries.len(), 2); + } } else { - assert!(false); + panic!("Root should be internal after split"); } } @@ -435,30 +493,27 @@ mod tests { tree.insert(i, vec![i as u8]).unwrap(); } - let mut current = Some(Rc::new(RefCell::new(tree.root))); - while let Some(node) = current.clone() { - let borrowed = node.borrow(); - match &*borrowed { + let mut current = tree.root.clone(); + loop { + let current_clone = current.clone(); + let node = current_clone.borrow(); + match &*node { Node::Internal(internal) => { - current = Some(Rc::clone(&internal.children[0])); + current = internal.children[0].clone(); } Node::Leaf(_) => break, } } - 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())))); - - while let Some(leaf_ref) = current_leaf { - let leaf_guard = leaf_ref.borrow(); - - 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 mut collected = Vec::new(); + let mut leaf_opt = Some(current); + while let Some(leaf_ref) = leaf_opt { + let leaf = leaf_ref.borrow(); + if let Node::Leaf(leaf_node) = &*leaf { + collected.extend(leaf_node.entries.iter().map(|(k, _)| **k)); + leaf_opt = leaf_node.next.as_ref().map(|n| n.clone()); + } else { + panic!("Expected leaf node"); } } @@ -478,7 +533,6 @@ mod tests { assert!(tree.file_number > 1); } - #[test] fn test_node_consistency_after_splits() { let tempdir = TempDir::new("consistency").unwrap(); @@ -488,14 +542,18 @@ mod tests { tree.insert(i, vec![i as u8]).unwrap(); } - if let Node::Internal(root) = &tree.root { - assert!(!root.keys.is_empty()); - assert!(root.children.len() >= 2); + let root = tree.root.borrow(); + if let Node::Internal(root_node) = &*root { + assert!(!root_node.keys.is_empty()); + assert!(root_node.children.len() >= 2); - for key in &root.keys { + for key in &root_node.keys { let key_val = **key; assert!(key_val > 0 && key_val < 10); } + } else { + panic!("Root should be internal node after splits"); } } } + From 1b96905a9f452db3d548706bc38a298d73c39857 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Sun, 4 May 2025 22:12:31 +0300 Subject: [PATCH 02/22] fix format --- src/bplus_tree.rs | 1 - 1 file changed, 1 deletion(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index 5ac18eb..08cd5aa 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -556,4 +556,3 @@ mod tests { } } } - From c8afa31697d5336019e7a6df96e2d77e7d99c758 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Sun, 4 May 2025 22:19:50 +0300 Subject: [PATCH 03/22] merge Node::insert/get functions into BPlus::insert/get functions --- src/bplus_tree.rs | 170 +++++++++++++++++++++------------------------- 1 file changed, 76 insertions(+), 94 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index 08cd5aa..e864f3e 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -175,90 +175,10 @@ impl BPlus { /// /// 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(); - let key = Rc::new(key); - - if let Some((new_node_link, new_key)) = - Node::insert(self.root.clone(), key, value_to_insert, self.t) - { - //let mut prev_root = self.root.clone(); - if let Node::Leaf(ref mut leaf) = &mut *(self.root.clone()).borrow_mut() { - leaf.next = Some(new_node_link.clone()); - } - 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(self.root.clone()); - 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 = Rc::new(RefCell::new(new_root)); - } - - Ok(()) - } - - #[allow(unused_variables)] - fn remove(&mut self, key: Rc) -> io::Result<()> { - unimplemented!() - } - - /// Gets value from a B+ tree by given key - pub fn get(&self, key: &K) -> io::Result> { - Node::get(self.root.clone(), key) - } -} - -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(); - - let new_leaf = Node::Leaf(Leaf { - entries: new_leaf_entries, - next: leaf.next.take(), - }); - - let new_leaf_link = Rc::new(RefCell::new(new_leaf.clone())); - 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); - - 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, - }); - - (Rc::new(RefCell::new(new_node)), middle_key) - } - } - } - - /// Inserts given value by given key - fn insert( - root: Link, - key: Rc, - value: ChunkHandler, - t: usize, - ) -> Option<(Link, Rc)> { + let value = self.get_chunk_handler(value).unwrap(); let mut path = Vec::new(); // Path to leaf - let mut current = root; + let mut current = self.root.clone(); let mut split_result; // Descent to the leaf @@ -272,8 +192,8 @@ impl Node { Err(pos) => leaf.entries.insert(pos, (key.clone(), value)), }; - split_result = if leaf.entries.len() == 2 * t { - Some(current_node.split(t)) + split_result = if leaf.entries.len() == 2 * self.t { + Some(current_node.split(self.t)) } else { None }; @@ -302,8 +222,8 @@ impl Node { internal.keys.insert(pos, median.clone()); internal.children.insert(pos + 1, new_node); - if internal.keys.len() == 2 * t - 1 { - split_result = Some(node.split(t)); + if internal.keys.len() == 2 * self.t - 1 { + split_result = Some(node.split(self.t)); } else { split_result = None; } @@ -313,17 +233,36 @@ impl Node { } } - split_result + // Splitting root if needed + if let Some((new_node_link, new_key)) = split_result { + //let mut prev_root = self.root.clone(); + if let Node::Leaf(ref mut leaf) = &mut *(self.root.clone()).borrow_mut() { + leaf.next = Some(new_node_link.clone()); + } + 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(self.root.clone()); + 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 = Rc::new(RefCell::new(new_root)); + } + + Ok(()) } - #[allow(unused_variables, dead_code)] - fn remove(&mut self, key: &K, t: usize) -> io::Result<()> { + #[allow(unused_variables)] + fn remove(&mut self, key: Rc) -> io::Result<()> { unimplemented!() } /// Gets value from a B+ tree by given key - fn get(root: Link, key: &K) -> io::Result> { - let mut current = root; + pub fn get(&self, key: &K) -> io::Result> { + let mut current = self.root.clone(); loop { let current_clone = current.clone(); @@ -352,11 +291,54 @@ impl Node { } } - drop(node); + drop(node); // Not sure is needed } } } +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(); + + let new_leaf = Node::Leaf(Leaf { + entries: new_leaf_entries, + next: leaf.next.take(), + }); + + let new_leaf_link = Rc::new(RefCell::new(new_leaf.clone())); + 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); + + 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, + }); + + (Rc::new(RefCell::new(new_node)), middle_key) + } + } + } + + #[allow(unused_variables, dead_code)] + fn remove(&mut self, key: &K, t: usize) -> io::Result<()> { + unimplemented!() + } +} + impl Database> for BPlus { fn insert(&mut self, key: K, value: DataContainer<()>) -> io::Result<()> { match value.extract() { @@ -366,11 +348,11 @@ impl Database> for BPlus< } fn get(&self, key: &K) -> io::Result> { - Node::get(self.root.clone(), key).map(DataContainer::from) + self.get(key).map(DataContainer::from) } fn contains(&self, key: &K) -> bool { - Node::get(self.root.clone(), key).is_ok() + self.get(key).is_ok() } } From 252434d9c59e7fec9f687e605d8b1039e8c5b1d8 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Wed, 21 May 2025 01:36:06 +0300 Subject: [PATCH 04/22] make insert/get async and threadsafe through latch crabbing, remove(for now) iterator and display trait realisation, most tests, implementation of chunkfs Database trait, etc --- Cargo.toml | 4 + src/bplus_tree.rs | 486 ++++++++++++++++++----------------------- tests/bplus_tests.rs | 4 +- tests/chunkfs_tests.rs | 4 +- 4 files changed, 216 insertions(+), 282 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index da763fc..2692094 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -8,8 +8,12 @@ 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"] } +env_logger = "0.11.8" +log = "0.4.27" diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index e864f3e..25d96dc 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -1,15 +1,20 @@ -use chunkfs::{Data, DataContainer, Database}; - use std::{ - cell::RefCell, - fmt::{self, Debug}, + collections::VecDeque, + fmt::Debug, fs::{create_dir_all, File}, io::{self, ErrorKind}, + mem, os::unix::fs::FileExt, path::PathBuf, rc::Rc, + sync::{ + atomic::{AtomicU64, AtomicUsize}, + Arc, + }, }; +use tokio::{self, sync::RwLock}; + const DEFAULT_MAX_FILE_SIZE: u64 = 2 << 20; extern crate chunkfs; @@ -43,7 +48,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 +65,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>, } @@ -81,59 +86,19 @@ pub struct BPlus { /// 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, -} - -impl fmt::Display for BPlus { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - BPlus::fmt_node(&self.root, 0, f) - } -} - -impl BPlus { - fn fmt_node(node: &Link, level: usize, f: &mut fmt::Formatter<'_>) -> fmt::Result { - let node_ref = node.borrow(); - - match &*node_ref { - Node::Internal(internal) => { - writeln!( - f, - "{}[Internal] keys: {:?}", - " ".repeat(level), - internal.keys.iter().map(|k| k.as_ref()).collect::>() - )?; - - for child in &internal.children { - Self::fmt_node(child, level + 1, f)?; - } - Ok(()) - } - Node::Leaf(leaf) => { - writeln!( - f, - "{}[Leaf] entries: {:?}, next: {}", - " ".repeat(level), - leaf.entries - .iter() - .map(|(k, v)| (k.as_ref(), v)) - .collect::>(), - leaf.next - .as_ref() - .map_or("None".into(), |n| format!("{:p}", Rc::as_ptr(n))) - ) - } - } - } + // Latch for root + latch: RwLock<()>, } #[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 @@ -142,49 +107,72 @@ impl BPlus { create_dir_all(&path)?; let current_file = File::create(path_to_file)?; Ok(Self { - root: Rc::new(RefCell::new(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); + *file_guard = File::create( + self.path.join( + self.file_number + .load(std::sync::atomic::Ordering::SeqCst) + .to_string(), + ), + ) + .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 key = Rc::new(key); - let value = self.get_chunk_handler(value).unwrap(); + pub async fn insert(&self, key: K, value: Vec) -> io::Result<()> { + let key = Arc::new(key); + let value = self.get_chunk_handler(value).await.unwrap(); let mut path = Vec::new(); // Path to leaf + let mut latch_guard = Some(self.latch.write()); let mut current = self.root.clone(); let mut split_result; + let mut guards = VecDeque::new(); // Descent to the leaf loop { - let current_clone = current.clone(); - let mut current_node = current_clone.borrow_mut(); + let mut current_node = current.write_owned().await; + if let Some(guard) = latch_guard { + drop(guard); + latch_guard = None; + }; match &mut *current_node { Node::Leaf(leaf) => { match leaf.entries.binary_search_by(|(k, _)| k.cmp(&key)) { @@ -192,11 +180,19 @@ impl BPlus { Err(pos) => leaf.entries.insert(pos, (key.clone(), value)), }; + println!("leaf.len = {}", leaf.entries.len()); + split_result = if leaf.entries.len() == 2 * self.t { Some(current_node.split(self.t)) } else { None }; + + // if path is empty, then current node is root + if path.is_empty() { + guards.push_back(current_node); + } + break; } Node::Internal(internal) => { @@ -205,77 +201,117 @@ impl BPlus { Err(pos) => pos, }; + // droping guards if nodes are not going to be changed + if internal.keys.len() != 2 * self.t - 2 { + while guards.len() > 1 { + drop(guards.pop_front().unwrap()); + } + } + let next_node = internal.children[pos].clone(); - path.push((current, pos)); + path.push(pos); current = next_node; } } + + guards.push_back(current_node); } // Going up to the root splitting nodes if needed - while let Some((parent, pos)) = path.pop() { + while let Some(pos) = path.pop() { if let Some((new_node, median)) = split_result.take() { - let mut node = parent.borrow_mut(); + 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 { - break; // No need to split the nodes anymore + drop(node); } } } - // Splitting root if needed - if let Some((new_node_link, new_key)) = split_result { - //let mut prev_root = self.root.clone(); - if let Node::Leaf(ref mut leaf) = &mut *(self.root.clone()).borrow_mut() { - leaf.next = Some(new_node_link.clone()); + // 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; + } + } + } } - 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(self.root.clone()); - 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 = Rc::new(RefCell::new(new_root)); } Ok(()) } - #[allow(unused_variables)] fn remove(&mut self, key: Rc) -> io::Result<()> { unimplemented!() } /// Gets value from a B+ tree by given key - pub fn get(&self, key: &K) -> io::Result> { + 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 current_clone = current.clone(); - let node = current_clone.borrow(); - + 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(_) => Err(ErrorKind::NotFound.into()), + Err(_) => { + drop(node); + Err(ErrorKind::NotFound.into()) + } }; } Node::Internal(internal) => { @@ -286,19 +322,21 @@ impl BPlus { current = match internal.children.get(pos) { Some(child) => child.clone(), - None => return Err(ErrorKind::NotFound.into()), + None => { + drop(node); + return Err(ErrorKind::NotFound.into()); + } }; } } - - drop(node); // Not sure is needed + 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) { + fn split(&mut self, t: usize) -> (Link, Arc) { match self { Node::Leaf(leaf) => { let mut new_leaf_entries = leaf.entries.split_off(t); @@ -310,7 +348,7 @@ impl Node { next: leaf.next.take(), }); - let new_leaf_link = Rc::new(RefCell::new(new_leaf.clone())); + let new_leaf_link = Arc::new(RwLock::new(new_leaf)); leaf.next = Some(new_leaf_link.clone()); (new_leaf_link, middle_key) @@ -328,7 +366,7 @@ impl Node { keys: new_node_keys, }); - (Rc::new(RefCell::new(new_node)), middle_key) + (Arc::new(RwLock::new(new_node)), middle_key) } } } @@ -339,202 +377,94 @@ impl Node { } } -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 get(&self, key: &K) -> io::Result> { - self.get(key).map(DataContainer::from) - } +#[cfg(test)] +mod tests { + use super::*; + use tempfile::TempDir; + use tokio::test; - fn contains(&self, key: &K) -> bool { - self.get(key).is_ok() + async 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) } -} -/// Iterator for B+ tree -pub struct BPlusIterator { - current_leaf: Option>>>, - current_index: usize, -} - -impl Iterator for BPlusIterator { - type Item = (K, Vec); + #[tokio::test(flavor = "multi_thread")] + async fn test_multiple_inserts() { + let (tree, _temp) = create_test_tree(2, "multiple_inserts").await; - 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; - } + for i in 1..=4 { + tree.insert(i, vec![i as u8]).await.unwrap(); } - } -} -impl IntoIterator for BPlus { - type Item = (K, Vec); - type IntoIter = BPlusIterator; - fn into_iter(self) -> Self::IntoIter { - let mut current = Some(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, - } - } - - BPlusIterator { - current_leaf: current, - current_index: 0, + for i in 1..=1 { + let result = tree.get(&i).await.unwrap(); + assert_eq!(result, vec![i as u8]); } } -} - -#[cfg(test)] -mod tests { - use tempdir::TempDir; - - use super::*; - #[test] - fn test_node_split_leaf() { - let tempdir = TempDir::new("split_leaf").unwrap(); - let mut tree = BPlus::new(2, tempdir.path().into()).unwrap(); - - for i in 1..5 { - tree.insert(i, vec![i as u8]).unwrap(); + #[tokio::test(flavor = "multi_thread")] + async fn test_concurrent_inserts() { + let (tree, _temp) = create_test_tree(2, "concurrent_inserts").await; + 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.unwrap(); + })); } - let root = tree.root.borrow(); - if let Node::Internal(root_node) = &*root { - assert_eq!(root_node.keys.len(), 1); - assert_eq!(root_node.children.len(), 2); - - let left_child = root_node.children[0].borrow(); - if let Node::Leaf(left_leaf) = &*left_child { - assert_eq!(left_leaf.entries.len(), 2); - } - - let right_child = root_node.children[1].borrow(); - if let Node::Leaf(right_leaf) = &*right_child { - assert_eq!(right_leaf.entries.len(), 2); - } - } else { - panic!("Root should be internal after split"); + for handle in handles { + handle.await.unwrap(); } - } - - #[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; - let data = vec![0; 200]; - tree.insert(1, data).unwrap(); - - let small_data = vec![0; 10]; - tree.insert(2, small_data).unwrap(); - - assert_eq!(tree.file_number, 1); - assert_eq!(tree.offset, 10); + let tree = tree.read().await; + for i in 0..50 { + 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(); + async fn test_root_split() { + let (tree, _temp) = create_test_tree(2, "root_split").await; - for i in 1..=9 { - tree.insert(i, vec![i as u8]).unwrap(); - } + tree.insert(1, vec![1]).await.unwrap(); + tree.insert(2, vec![2]).await.unwrap(); + tree.insert(3, vec![3]).await.unwrap(); + tree.insert(4, vec![4]).await.unwrap(); - let mut current = tree.root.clone(); - loop { - let current_clone = current.clone(); - let node = current_clone.borrow(); - match &*node { - Node::Internal(internal) => { - current = internal.children[0].clone(); - } - Node::Leaf(_) => break, - } - } - - let mut collected = Vec::new(); - let mut leaf_opt = Some(current); - while let Some(leaf_ref) = leaf_opt { - let leaf = leaf_ref.borrow(); - if let Node::Leaf(leaf_node) = &*leaf { - collected.extend(leaf_node.entries.iter().map(|(k, _)| **k)); - leaf_opt = leaf_node.next.as_ref().map(|n| n.clone()); - } else { - panic!("Expected leaf node"); + 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; - - for i in 0..500 { - tree.insert(i, vec![i as u8; 100]).unwrap(); - } - - assert!(tree.file_number > 1); - } - #[test] - fn test_node_consistency_after_splits() { - let tempdir = TempDir::new("consistency").unwrap(); - let mut tree = BPlus::new(2, tempdir.path().into()).unwrap(); - - for i in 0..10 { - tree.insert(i, vec![i as u8]).unwrap(); - } - - let root = tree.root.borrow(); - if let Node::Internal(root_node) = &*root { - assert!(!root_node.keys.is_empty()); - assert!(root_node.children.len() >= 2); - - for key in &root_node.keys { - let key_val = **key; - assert!(key_val > 0 && key_val < 10); - } - } else { - panic!("Root should be internal node after splits"); - } + 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; + + let large_data = vec![7; 150]; + tree.insert(1, large_data.clone()).await.unwrap(); + + let result = tree.get(&1).await.unwrap(); + assert_eq!(result, large_data); + tree.insert(2, large_data.clone()).await.unwrap(); + 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" + ); } } diff --git a/tests/bplus_tests.rs b/tests/bplus_tests.rs index 1a6609e..d7a81e8 100644 --- a/tests/bplus_tests.rs +++ b/tests/bplus_tests.rs @@ -1,4 +1,4 @@ -extern crate chunkfs; +/*extern crate chunkfs; use std::collections::HashMap; use std::path::PathBuf; @@ -317,4 +317,4 @@ fn test_custom_data_types() { assert_eq!(tree.get(&key1).unwrap(), b"data1"); assert_eq!(tree.get(&key2).unwrap(), b"data2"); -} +}*/ diff --git a/tests/chunkfs_tests.rs b/tests/chunkfs_tests.rs index 23c4b89..20f8164 100644 --- a/tests/chunkfs_tests.rs +++ b/tests/chunkfs_tests.rs @@ -1,4 +1,4 @@ -extern crate chunkfs; +/*extern crate chunkfs; use std::collections::HashMap; use std::io; @@ -218,4 +218,4 @@ fn write_from_stream_buf_reader() { let read = fs.read_file_complete(&ro_fh).unwrap(); assert_eq!(read.len(), MB); assert_eq!(read, [1; MB]); -} +}*/ From 4abec6f038a48aa30d322534fef7ddb54975f20a Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Thu, 29 May 2025 12:55:16 +0300 Subject: [PATCH 05/22] add BPlusStorage struct and Database trait implementation for it --- src/bplus_tree.rs | 93 +++++++++++++++++++++++++++++++++++++----- tests/chunkfs_tests.rs | 42 +++++++++++++------ 2 files changed, 112 insertions(+), 23 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index 25d96dc..40e1a69 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -1,5 +1,5 @@ use std::{ - collections::VecDeque, + collections::{HashSet, VecDeque}, fmt::Debug, fs::{create_dir_all, File}, io::{self, ErrorKind}, @@ -11,9 +11,15 @@ use std::{ atomic::{AtomicU64, AtomicUsize}, Arc, }, + thread, time, }; -use tokio::{self, sync::RwLock}; +use chunkfs::{Data, DataContainer, Database}; +use tokio::{ + self, + runtime::Runtime, + sync::{Mutex, RwLock}, +}; const DEFAULT_MAX_FILE_SIZE: u64 = 2 << 20; @@ -97,6 +103,72 @@ pub struct BPlus { 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 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 Database, DataContainer<()>> for BPlusStorage { + /// Inserts given value by given key in the B+ tree + fn insert(&mut self, key: Vec, 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(); + + self.runtime.spawn(async move { + set_clone.lock().await.insert(key.clone()); + tree.insert(key.clone(), value).await.unwrap(); + set_clone.lock().await.remove(&key); + }); + Ok(()) + } + + /// Gets value by given key from B+ tree + fn get(&self, key: &Vec) -> 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().await.contains(key) { + thread::sleep(time::Duration::from_millis(10)); + } + tree.get(key).await.unwrap() + }) + .into()) + } + + /// Returns whether key is contained in the B+ tree or not + fn contains(&self, key: &Vec) -> bool { + self.get(key).is_ok() + } +} + #[allow(dead_code)] impl BPlus { /// Creates new instance of B+ tree with given t and path @@ -106,6 +178,7 @@ impl BPlus { let path_to_file = path.join("0"); create_dir_all(&path)?; let current_file = File::create(path_to_file)?; + Ok(Self { root: Arc::new(RwLock::new(Node::Leaf(Leaf::default()))), t, @@ -383,29 +456,29 @@ mod tests { use tempfile::TempDir; use tokio::test; - async fn create_test_tree(t: usize, name: &str) -> (BPlus, TempDir) { + 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) } - #[tokio::test(flavor = "multi_thread")] + #[test] async fn test_multiple_inserts() { - let (tree, _temp) = create_test_tree(2, "multiple_inserts").await; + let (tree, _temp) = create_test_tree(2, "multiple_inserts"); for i in 1..=4 { - tree.insert(i, vec![i as u8]).await.unwrap(); + tree.insert(i, vec![i as u8]).await.unwrap(); // let _ = handle.spawn(async move {tree.insert(i, vec![i as u8]).await}).await.unwrap(); } - for i in 1..=1 { + for i in 1..=4 { let result = tree.get(&i).await.unwrap(); assert_eq!(result, vec![i as u8]); } } - #[tokio::test(flavor = "multi_thread")] + #[test] async fn test_concurrent_inserts() { - let (tree, _temp) = create_test_tree(2, "concurrent_inserts").await; + let (tree, _temp) = create_test_tree(2, "concurrent_inserts"); let tree = Arc::new(tokio::sync::RwLock::new(tree)); let mut handles = vec![]; @@ -430,7 +503,7 @@ mod tests { #[test] async fn test_root_split() { - let (tree, _temp) = create_test_tree(2, "root_split").await; + let (tree, _temp) = create_test_tree(2, "root_split"); tree.insert(1, vec![1]).await.unwrap(); tree.insert(2, vec![2]).await.unwrap(); diff --git a/tests/chunkfs_tests.rs b/tests/chunkfs_tests.rs index 20f8164..2795f7c 100644 --- a/tests/chunkfs_tests.rs +++ b/tests/chunkfs_tests.rs @@ -1,4 +1,4 @@ -/*extern crate chunkfs; +extern crate chunkfs; use std::collections::HashMap; use std::io; @@ -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(); @@ -218,4 +234,4 @@ fn write_from_stream_buf_reader() { let read = fs.read_file_complete(&ro_fh).unwrap(); assert_eq!(read.len(), MB); assert_eq!(read, [1; MB]); -}*/ +} From e0fdcbb2b96f7ac438248cee992288ce6f21bed8 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Thu, 29 May 2025 13:18:08 +0300 Subject: [PATCH 06/22] Add pending_inserts field to the BPlusStorage to fix tests, remove debug output --- src/bplus_tree.rs | 16 +++++++++++++--- 1 file changed, 13 insertions(+), 3 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index 40e1a69..c443ebc 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -111,6 +111,7 @@ pub struct BPlusStorage { runtime: Runtime, /// Currently inserting keys keys_set: Arc>>>, + pending_inserts: Arc, } impl BPlusStorage { @@ -124,6 +125,7 @@ impl BPlusStorage { tree: Arc::new(tree), runtime, keys_set: Arc::new(Mutex::new(HashSet::new())), + pending_inserts: Arc::new(0.into()), }) } } @@ -139,9 +141,12 @@ impl Database, DataContainer<()>> for BPlusStorage { }; let set_clone = self.keys_set.clone(); + let pending_inserts_clone = self.pending_inserts.clone(); + pending_inserts_clone.fetch_add(1, std::sync::atomic::Ordering::SeqCst); self.runtime.spawn(async move { set_clone.lock().await.insert(key.clone()); + pending_inserts_clone.fetch_sub(1, std::sync::atomic::Ordering::SeqCst); tree.insert(key.clone(), value).await.unwrap(); set_clone.lock().await.remove(&key); }); @@ -152,6 +157,13 @@ impl Database, DataContainer<()>> for BPlusStorage { fn get(&self, key: &Vec) -> io::Result> { let tree = self.tree.clone(); let set_clone = self.keys_set.clone(); + if self + .pending_inserts + .load(std::sync::atomic::Ordering::SeqCst) + > 0 + { + thread::sleep(time::Duration::from_millis(10)); + } Ok(self .runtime .block_on(async move { @@ -253,8 +265,6 @@ impl BPlus { Err(pos) => leaf.entries.insert(pos, (key.clone(), value)), }; - println!("leaf.len = {}", leaf.entries.len()); - split_result = if leaf.entries.len() == 2 * self.t { Some(current_node.split(self.t)) } else { @@ -467,7 +477,7 @@ mod tests { let (tree, _temp) = create_test_tree(2, "multiple_inserts"); for i in 1..=4 { - tree.insert(i, vec![i as u8]).await.unwrap(); // let _ = handle.spawn(async move {tree.insert(i, vec![i as u8]).await}).await.unwrap(); + tree.insert(i, vec![i as u8]).await.unwrap(); } for i in 1..=4 { From 6664cf0afef3e94b2b0ea4fcbd88b0265da1c66c Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Thu, 29 May 2025 13:35:11 +0300 Subject: [PATCH 07/22] remove pending_inserts, change tokio mutex to std mutex --- src/bplus_tree.rs | 27 ++++++--------------------- 1 file changed, 6 insertions(+), 21 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index c443ebc..875acdc 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -9,17 +9,13 @@ use std::{ rc::Rc, sync::{ atomic::{AtomicU64, AtomicUsize}, - Arc, + Arc, Mutex, }, thread, time, }; use chunkfs::{Data, DataContainer, Database}; -use tokio::{ - self, - runtime::Runtime, - sync::{Mutex, RwLock}, -}; +use tokio::{self, runtime::Runtime, sync::RwLock}; const DEFAULT_MAX_FILE_SIZE: u64 = 2 << 20; @@ -111,7 +107,6 @@ pub struct BPlusStorage { runtime: Runtime, /// Currently inserting keys keys_set: Arc>>>, - pending_inserts: Arc, } impl BPlusStorage { @@ -125,7 +120,6 @@ impl BPlusStorage { tree: Arc::new(tree), runtime, keys_set: Arc::new(Mutex::new(HashSet::new())), - pending_inserts: Arc::new(0.into()), }) } } @@ -141,14 +135,11 @@ impl Database, DataContainer<()>> for BPlusStorage { }; let set_clone = self.keys_set.clone(); - let pending_inserts_clone = self.pending_inserts.clone(); - pending_inserts_clone.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + set_clone.lock().unwrap().insert(key.clone()); self.runtime.spawn(async move { - set_clone.lock().await.insert(key.clone()); - pending_inserts_clone.fetch_sub(1, std::sync::atomic::Ordering::SeqCst); tree.insert(key.clone(), value).await.unwrap(); - set_clone.lock().await.remove(&key); + set_clone.lock().unwrap().remove(&key); }); Ok(()) } @@ -157,17 +148,11 @@ impl Database, DataContainer<()>> for BPlusStorage { fn get(&self, key: &Vec) -> io::Result> { let tree = self.tree.clone(); let set_clone = self.keys_set.clone(); - if self - .pending_inserts - .load(std::sync::atomic::Ordering::SeqCst) - > 0 - { - thread::sleep(time::Duration::from_millis(10)); - } + Ok(self .runtime .block_on(async move { - while set_clone.lock().await.contains(key) { + while set_clone.lock().unwrap().contains(key) { thread::sleep(time::Duration::from_millis(10)); } tree.get(key).await.unwrap() From 17d80d7723d78db33f00bbf426f963659d9b3636 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Mon, 2 Jun 2025 21:47:03 +0300 Subject: [PATCH 08/22] change test so they work correct with async functions, add concurrent stress test --- src/bplus_tree.rs | 9 +- tests/bplus_tests.rs | 316 +++++++++++++++++++------------------------ 2 files changed, 140 insertions(+), 185 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index 875acdc..d465dcc 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -449,7 +449,6 @@ impl Node { mod tests { use super::*; use tempfile::TempDir; - use tokio::test; fn create_test_tree(t: usize, name: &str) -> (BPlus, TempDir) { let temp_dir = TempDir::with_prefix(name).unwrap(); @@ -457,7 +456,7 @@ mod tests { (tree, temp_dir) } - #[test] + #[tokio::test(flavor = "multi_thread")] async fn test_multiple_inserts() { let (tree, _temp) = create_test_tree(2, "multiple_inserts"); @@ -471,7 +470,7 @@ mod tests { } } - #[test] + #[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)); @@ -496,7 +495,7 @@ mod tests { } } - #[test] + #[tokio::test(flavor = "multi_thread")] async fn test_root_split() { let (tree, _temp) = create_test_tree(2, "root_split"); @@ -515,7 +514,7 @@ mod tests { } } - #[test] + #[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(); diff --git a/tests/bplus_tests.rs b/tests/bplus_tests.rs index d7a81e8..c03a367 100644 --- a/tests/bplus_tests.rs +++ b/tests/bplus_tests.rs @@ -1,115 +1,112 @@ -/*extern crate chunkfs; +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.unwrap(); + 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.unwrap(); + tree.insert(1, vec![42]).await.unwrap(); - 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.unwrap(); } 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.unwrap(); } 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.unwrap(); } 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.unwrap(); 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.unwrap(); } 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.unwrap(); } } 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 { let key: usize = rand::random::() % 10000; @@ -117,204 +114,163 @@ fn test_same_keys_inserted() { } for key in keys.clone() { - tree.insert(key, vec![key as u8]).unwrap(); + tree.insert(key, vec![key as u8]).await.unwrap(); } 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.unwrap(); 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.unwrap(); } } -#[test] -fn test_iterator() { - let tempdir = TempDir::new("iterator_test").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(); - } - - 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); -} -#[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.unwrap(); + 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.unwrap(); } for i in 1..100 { - assert_eq!(tree.get(&i).unwrap(), vec![i as u8]); + assert_eq!(tree.get(&i).await.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])); - } - - 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.unwrap(); } - 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.unwrap(); + tree.insert(42, vec![2]).await.unwrap(); } - 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("apple".to_string(), b"fruit".to_vec()) + .await + .unwrap(); tree.insert("banana".to_string(), b"yellow".to_vec()) + .await .unwrap(); - 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.unwrap(); } 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]); } } -#[test] -fn test_find_nonexistent_after_splits() { - let tempdir = TempDir::new("nonexistent").unwrap(); - let mut tree = BPlus::new(2, tempdir.path().into()).unwrap(); +#[tokio::test(flavor = "multi_thread", worker_threads = 16)] +async fn test_stress_1m_entries_concurrent() { + use std::sync::Arc; + use tokio::task; - for i in 0..1000 { - tree.insert(i, vec![1]).unwrap(); + 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.unwrap(); + } + })); } - assert!(tree.get(&1001).is_err()); -} - -#[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); + barrier.wait().await; - for &num in &nums { - tree.insert(num, vec![num as u8]).unwrap(); + for handle in insert_handles { + handle.await.unwrap(); } - let mut sorted = nums.clone(); - sorted.sort_unstable(); - let result: Vec<_> = tree.into_iter().map(|(k, _)| k).collect(); - assert_eq!(result, sorted); -} - -#[test] -fn test_custom_data_types() { - #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Default)] - struct ComplexKey { - id: u64, - name: String, + 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); + } + })); } - let tempdir = TempDir::new("custom_type").unwrap(); - let mut tree = BPlus::new(2, tempdir.path().into()).unwrap(); + for handle in verify_handles { + handle.await.unwrap(); + } +} - let key1 = ComplexKey { - id: 1, - name: "A".to_string(), - }; - let key2 = ComplexKey { - id: 2, - name: "B".to_string(), - }; +#[tokio::test(flavor = "multi_thread")] +async fn test_find_nonexistent_after_splits() { + let tempdir = TempDir::new("nonexistent").unwrap(); + let tree = BPlus::new(2, tempdir.path().into()).unwrap(); - tree.insert(key1.clone(), b"data1".to_vec()).unwrap(); - tree.insert(key2.clone(), b"data2".to_vec()).unwrap(); + for i in 0..1000 { + tree.insert(i, vec![1]).await.unwrap(); + } - assert_eq!(tree.get(&key1).unwrap(), b"data1"); - assert_eq!(tree.get(&key2).unwrap(), b"data2"); -}*/ + assert!(tree.get(&1001).await.is_err()); +} From 93c35723a0636fa893d0c7c1a6d846a866f7d719 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Fri, 13 Jun 2025 20:26:37 +0300 Subject: [PATCH 09/22] add droping guards in leaf nodes in insert functions --- src/bplus_tree.rs | 19 +++++++++++-------- 1 file changed, 11 insertions(+), 8 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index d465dcc..dee175f 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -100,16 +100,16 @@ pub struct BPlus { } /// Wrapper for BPlusTree with sync functions with async runtime -pub struct BPlusStorage { +pub struct BPlusStorage { /// BPlusTree - tree: Arc>>, + tree: Arc>, /// Async tokio runtime for operations runtime: Runtime, /// Currently inserting keys - keys_set: Arc>>>, + keys_set: Arc>>, } -impl BPlusStorage { +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 @@ -124,9 +124,9 @@ impl BPlusStorage { } } -impl Database, DataContainer<()>> for BPlusStorage { +impl Database> for BPlusStorage { /// Inserts given value by given key in the B+ tree - fn insert(&mut self, key: Vec, value: DataContainer<()>) -> io::Result<()> { + fn insert(&mut self, key: K, value: DataContainer<()>) -> io::Result<()> { let tree = self.tree.clone(); let value = match value.extract() { @@ -145,7 +145,7 @@ impl Database, DataContainer<()>> for BPlusStorage { } /// Gets value by given key from B+ tree - fn get(&self, key: &Vec) -> io::Result> { + fn get(&self, key: &K) -> io::Result> { let tree = self.tree.clone(); let set_clone = self.keys_set.clone(); @@ -161,7 +161,7 @@ impl Database, DataContainer<()>> for BPlusStorage { } /// Returns whether key is contained in the B+ tree or not - fn contains(&self, key: &Vec) -> bool { + fn contains(&self, key: &K) -> bool { self.get(key).is_ok() } } @@ -253,6 +253,9 @@ impl BPlus { 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 }; From e0f13ef4394d7e1dff901228ce8e388669a8d778 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Fri, 13 Jun 2025 20:27:11 +0300 Subject: [PATCH 10/22] fix formatting --- src/bplus_tree.rs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index dee175f..a171046 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -124,7 +124,9 @@ impl BPlusStorage { } } -impl Database> for BPlusStorage { +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(); From e62d3136e2435bc798222fa65a426d3ca25ea7c4 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Tue, 17 Jun 2025 00:36:00 +0300 Subject: [PATCH 11/22] Add save and load functions; add serializable tree and node structures, tests, some other minor changes --- src/bplus_tree.rs | 315 +++++++++++++++++++++++++++++++++++++++++-- tests/bplus_tests.rs | 20 ++- 2 files changed, 325 insertions(+), 10 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index a171046..23ea608 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -2,18 +2,22 @@ use std::{ 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}, + 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}; @@ -21,8 +25,119 @@ const DEFAULT_MAX_FILE_SIZE: u64 = 2 << 20; extern crate chunkfs; +#[derive(Serialize, Deserialize)] +struct SerializableBPlus { + t: usize, + path: PathBuf, + file_number: usize, + offset: u64, + max_file_size: u64, + root: SerializableNode, +} + +#[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 { + 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] + 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 Deserialize<'de>> + SerializableBPlus +{ + 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, @@ -109,7 +224,9 @@ pub struct BPlusStorage { keys_set: Arc>>, } -impl BPlusStorage { +impl Deserialize<'de> + Sync + Send> + 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 @@ -124,8 +241,18 @@ impl BPlusStorage { } } -impl - Database> for BPlusStorage +impl< + K: Clone + + Ord + + std::hash::Hash + + Debug + + Default + + Send + + Sync + + Serialize + + for<'de> Deserialize<'de> + + 'static, + > Database> for BPlusStorage { /// Inserts given value by given key in the B+ tree fn insert(&mut self, key: K, value: DataContainer<()>) -> io::Result<()> { @@ -169,7 +296,10 @@ impl } #[allow(dead_code)] -impl BPlus { +impl< + K: Default + Ord + Clone + Debug + Sized + Serialize + for<'de> Deserialize<'de> + Sync + Send, + > 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 @@ -264,6 +394,8 @@ impl BPlus { // if path is empty, then current node is root if path.is_empty() { guards.push_back(current_node); + } else { + drop(current_node); } break; @@ -276,7 +408,7 @@ impl BPlus { // droping guards if nodes are not going to be changed if internal.keys.len() != 2 * self.t - 2 { - while guards.len() > 1 { + while !guards.is_empty() { drop(guards.pop_front().unwrap()); } } @@ -347,12 +479,18 @@ impl BPlus { *node = new_root; } } + drop(node); } } } + for guard in guards { + drop(guard); + } + Ok(()) } + #[allow(unused_variables)] fn remove(&mut self, key: Rc) -> io::Result<()> { unimplemented!() @@ -405,6 +543,92 @@ impl BPlus { prev_guard = Some(node); } } + + 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; + } + + 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!(), + } + } + }) + .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); + } + } + } + + async fn collect_leaves(&self) -> Vec>>> { + let mut leaves = Vec::new(); + let mut queue = VecDeque::new(); + queue.push_back(self.root.clone()); + + 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()); + } + } + Node::Leaf(_) => { + leaves.push(node.clone()); + } + } + } + + leaves + } + + fn open_current_file(path: &Path, number: usize) -> io::Result>> { + Ok(Arc::new(RwLock::new( + File::open(path.join(number.to_string())).unwrap(), + ))) + } + + /// 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(|e| io::Error::new(io::ErrorKind::Other, e)) + } + + /// 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(|e| io::Error::new(io::ErrorKind::Other, e))?; + + Ok(serializable.deserialize().await) + } } impl Node { @@ -539,4 +763,77 @@ mod tests { "Should create multiple files" ); } + + #[tokio::test] + async fn test_save_load_empty_tree() { + let tempdir = TempDir::new().unwrap(); + let tree_path = tempdir.path().join("empty_tree.bin"); + + let tree = BPlus::::new(2, tempdir.path().into()).unwrap(); + + tree.save(&tree_path).await.unwrap(); + + 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()); + } + + #[tokio::test] + async fn test_save_load_small_tree() { + let tempdir = TempDir::new().unwrap(); + let tree_path = tempdir.path().join("small_tree.bin"); + + let tree = BPlus::::new(2, tempdir.path().into()).unwrap(); + tree.insert(10, vec![1, 2, 3]).await.unwrap(); + tree.insert(20, vec![4, 5, 6]).await.unwrap(); + tree.insert(5, vec![0]).await.unwrap(); + + tree.save(&tree_path).await.unwrap(); + + 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()); + } + + #[tokio::test] + async fn test_save_load_large_tree() { + let tempdir = TempDir::with_prefix("large_load_save").unwrap(); + let tree_path = tempdir.path().join("large_tree.bin"); + let mut tree = BPlus::::new(2, tempdir.path().into()).unwrap(); + tree.max_file_size = 100; + + for i in 0..100000 { + tree.insert(i, vec![(i % 256) as u8; 1]).await.unwrap(); + } + tree.save(&tree_path).await.unwrap(); + + let loaded_tree = BPlus::::load(&tree_path).await.unwrap(); + + use rand::seq::SliceRandom; + use rand::thread_rng; + + let mut rng = thread_rng(); + let mut keys: Vec = (0..10_000).collect(); + keys.shuffle(&mut rng); + + for key in keys.iter().take(100) { + let expected = vec![(*key % 256) as u8; 1]; + assert_eq!(loaded_tree.get(key).await.unwrap(), expected); + } + + assert!(loaded_tree.get(&100_000).await.is_err()); + } } diff --git a/tests/bplus_tests.rs b/tests/bplus_tests.rs index c03a367..47e1e62 100644 --- a/tests/bplus_tests.rs +++ b/tests/bplus_tests.rs @@ -108,7 +108,7 @@ async fn test_same_keys_inserted() { let tempdir = TempDir::new("10").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); } @@ -129,6 +129,24 @@ async fn test_same_keys_inserted() { } } +#[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 0..10000 { + tree.insert(i, vec![i as u8]).await.unwrap(); + } + + for i in 0..10000 { + tree.insert(i, vec![i as u8]).await.unwrap(); + } + + for key in 1..10000 { + assert_eq!(vec![key as u8], tree.get(&key).await.unwrap()); + } +} + #[tokio::test(flavor = "multi_thread")] async fn test_empty_tree() { let tempdir = TempDir::new("empty").unwrap(); From c5ca9a2cdd0272dd2e6879e8541e3743a18d0062 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Tue, 17 Jun 2025 00:40:23 +0300 Subject: [PATCH 12/22] add futures, bincode, serde --- Cargo.toml | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/Cargo.toml b/Cargo.toml index 2692094..0093a2f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -17,3 +17,7 @@ rand = "0.8.5" tokio = { version = "1.44.2", features = ["full"] } env_logger = "0.11.8" log = "0.4.27" +serde = { version = "1.0", features = ["derive"] } +bincode = "1.3" +async-recursion = "1.1.1" +futures = "0.3.31" From 010c93f34842d989116ad22f22446d3966b86a4b Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Tue, 17 Jun 2025 00:49:12 +0300 Subject: [PATCH 13/22] fix clippy warnings --- src/bplus_tree.rs | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index 23ea608..06e734c 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -616,16 +616,15 @@ impl< let serializable = self.serialize().await; let file = File::create(path)?; let writer = BufWriter::new(file); - bincode::serialize_into(writer, &serializable) - .map_err(|e| io::Error::new(io::ErrorKind::Other, e)) + bincode::serialize_into(writer, &serializable).map_err(io::Error::other) } /// 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(|e| io::Error::new(io::ErrorKind::Other, e))?; + let serializable: SerializableBPlus = + bincode::deserialize_from(reader).map_err(io::Error::other)?; Ok(serializable.deserialize().await) } From fdb05fe0adac530f79812ea323b34f26cfe6b28d Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Tue, 17 Jun 2025 01:44:56 +0300 Subject: [PATCH 14/22] add optimistic_insert function and change insert function for optimistic latch crabbing --- src/bplus_tree.rs | 79 ++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 78 insertions(+), 1 deletion(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index 06e734c..04c1758 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -360,10 +360,17 @@ impl< /// /// Returns Err(_) if file could not be created pub async fn insert(&self, key: K, value: Vec) -> io::Result<()> { - let key = Arc::new(key); let value = self.get_chunk_handler(value).await.unwrap(); let mut path = Vec::new(); // Path to leaf + if self + .optimistic_insert(key.clone(), value.clone()) + .await + .is_ok() + { + return Ok(()); + } 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(); @@ -544,6 +551,76 @@ impl< } } + /// 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 mut prev_guard = None; + loop { + let node = current.read_owned().await; + if let Some(guard) = latch_guard { + drop(guard); + latch_guard = None; + if let Node::Leaf(_) = &*node { + return Err(()); + } + } + + if let Node::Leaf(_) = node.clone() { + break; + } + + if prev_guard.is_some() { + drop(prev_guard); + } + + match &*node { + Node::Leaf(_leaf) => { + unreachable!() + } + Node::Internal(internal) => { + let pos = match internal.keys.binary_search(&key) { + Ok(pos) => pos + 1, + Err(pos) => pos, + }; + + let next_node = internal.children[pos].clone(); + + current = next_node; + } + } + prev_guard = Some(node); + } + + let current_node_guard = if let Node::Internal(internal) = &*prev_guard.take().unwrap() { + let pos = match internal.keys.binary_search(&key) { + Ok(pos) => pos + 1, + Err(pos) => pos, + }; + internal.clone().children[pos].clone() + } else { + unreachable!(); + }; + let mut current_node = current_node_guard.write().await; + if let Node::Leaf(leaf) = &mut *current_node { + if leaf.entries.len() == 2 * self.t - 1 { + return Err(()); + } + 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)), + }; + } + + Ok(()) + } + 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 From a722df3d876a6d277ee5f4c217d0b46bedb75fc5 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Tue, 17 Jun 2025 01:56:46 +0300 Subject: [PATCH 15/22] fix formatting --- src/bplus_tree.rs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index 04c1758..3e8c15c 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -362,6 +362,7 @@ impl< pub async fn insert(&self, key: K, value: Vec) -> io::Result<()> { 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 @@ -608,6 +609,7 @@ impl< unreachable!(); }; let mut current_node = current_node_guard.write().await; + if let Node::Leaf(leaf) = &mut *current_node { if leaf.entries.len() == 2 * self.t - 1 { return Err(()); From 9f3e0aaa05762d7f450e736e7bd758d971467105 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Tue, 17 Jun 2025 14:04:14 +0300 Subject: [PATCH 16/22] optimize optimistic_insert function --- src/bplus_tree.rs | 102 ++++++++++++++++++++-------------------------- 1 file changed, 45 insertions(+), 57 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index 3e8c15c..6878636 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -224,9 +224,7 @@ pub struct BPlusStorage { keys_set: Arc>>, } -impl Deserialize<'de> + Sync + Send> - BPlusStorage -{ +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 @@ -241,18 +239,8 @@ impl Deserialize<'de> + } } -impl< - K: Clone - + Ord - + std::hash::Hash - + Debug - + Default - + Send - + Sync - + Serialize - + for<'de> Deserialize<'de> - + 'static, - > Database> for BPlusStorage +impl + Database> for BPlusStorage { /// Inserts given value by given key in the B+ tree fn insert(&mut self, key: K, value: DataContainer<()>) -> io::Result<()> { @@ -296,10 +284,7 @@ impl< } #[allow(dead_code)] -impl< - K: Default + Ord + Clone + Debug + Sized + Serialize + for<'de> Deserialize<'de> + Sync + Send, - > 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 @@ -563,66 +548,69 @@ impl< let key = Arc::new(key); let mut prev_guard = None; + let mut last_child_index = None; + loop { let node = current.read_owned().await; - if let Some(guard) = latch_guard { + + if let Some(guard) = latch_guard.take() { drop(guard); - latch_guard = None; - if let Node::Leaf(_) = &*node { + if matches!(&*node, Node::Leaf(_)) { return Err(()); } } - if let Node::Leaf(_) = node.clone() { + if matches!(&*node, Node::Leaf(_)) { break; } - if prev_guard.is_some() { - drop(prev_guard); - } - - match &*node { - Node::Leaf(_leaf) => { - unreachable!() - } - Node::Internal(internal) => { - let pos = match internal.keys.binary_search(&key) { - Ok(pos) => pos + 1, - Err(pos) => pos, - }; - - let next_node = internal.children[pos].clone(); + prev_guard = Some(node); - current = next_node; - } + 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!(); } - prev_guard = Some(node); } - let current_node_guard = if let Node::Internal(internal) = &*prev_guard.take().unwrap() { - let pos = match internal.keys.binary_search(&key) { - Ok(pos) => pos + 1, - Err(pos) => pos, - }; - internal.clone().children[pos].clone() - } else { - unreachable!(); + let leaf_lock = match prev_guard.take() { + Some(parent) => { + let pos = last_child_index.unwrap(); + if let Node::Internal(internal) = &*parent { + internal.children[pos].clone() + } else { + unreachable!(); + } + } + None => unreachable!(), }; - let mut current_node = current_node_guard.write().await; - if let Node::Leaf(leaf) = &mut *current_node { - if leaf.entries.len() == 2 * self.t - 1 { + let mut leaf = leaf_lock.write().await; + if let Node::Leaf(leaf_node) = &mut *leaf { + if leaf_node.entries.len() == 2 * self.t - 1 { return Err(()); } - 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)), + + 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(()) + } else { + unreachable!() } - - Ok(()) } +} +impl< + K: Default + Ord + Clone + Debug + Sized + Serialize + for<'de> Deserialize<'de> + Sync + Send, + > BPlus +{ 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 From b565de15a8255fd592340144cdfb1d321783647c Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Tue, 17 Jun 2025 20:12:29 +0300 Subject: [PATCH 17/22] add comments, reorginize tests, other small changes --- src/bplus_tree.rs | 77 ++++++++---------------------------- tests/bplus_tests.rs | 92 ++++++++++++++++++++++++++++++-------------- 2 files changed, 80 insertions(+), 89 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index 6878636..93e2e0b 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -25,6 +25,7 @@ const DEFAULT_MAX_FILE_SIZE: u64 = 2 << 20; extern crate chunkfs; +/// Easily serializable version of BPlusTree #[derive(Serialize, Deserialize)] struct SerializableBPlus { t: usize, @@ -35,6 +36,7 @@ struct SerializableBPlus { root: SerializableNode, } +/// Easily serializable version of BPlusTree Node #[derive(Serialize, Deserialize)] enum SerializableNode { Internal(SerializableInternalNode), @@ -255,7 +257,7 @@ impl set_clone.lock().unwrap().insert(key.clone()); self.runtime.spawn(async move { - tree.insert(key.clone(), value).await.unwrap(); + tree.insert(key.clone(), value).await; set_clone.lock().unwrap().remove(&key); }); Ok(()) @@ -344,7 +346,7 @@ impl BPlus { /// Inserts given value by given key in the B+ tree /// /// Returns Err(_) if file could not be created - pub async fn insert(&self, key: K, value: Vec) -> io::Result<()> { + 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() @@ -353,7 +355,7 @@ impl BPlus { .await .is_ok() { - return Ok(()); + return; } let mut latch_guard = Some(self.latch.write()); let key = Arc::new(key); @@ -481,7 +483,7 @@ impl BPlus { drop(guard); } - Ok(()) + // Ok(()) } #[allow(unused_variables)] @@ -611,6 +613,7 @@ impl< K: Default + Ord + Clone + Debug + Sized + Serialize + for<'de> Deserialize<'de> + Sync + Send, > 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 @@ -649,6 +652,7 @@ impl< } } + /// Collects all leaves from BPlusTree async fn collect_leaves(&self) -> Vec>>> { let mut leaves = Vec::new(); let mut queue = VecDeque::new(); @@ -756,7 +760,7 @@ mod tests { let (tree, _temp) = create_test_tree(2, "multiple_inserts"); for i in 1..=4 { - tree.insert(i, vec![i as u8]).await.unwrap(); + tree.insert(i, vec![i as u8]).await; } for i in 1..=4 { @@ -775,7 +779,7 @@ mod tests { let tree = tree.clone(); handles.push(tokio::spawn(async move { let tree = tree.write().await; - tree.insert(i, vec![i as u8]).await.unwrap(); + tree.insert(i, vec![i as u8]).await; })); } @@ -794,10 +798,10 @@ mod tests { async fn test_root_split() { let (tree, _temp) = create_test_tree(2, "root_split"); - tree.insert(1, vec![1]).await.unwrap(); - tree.insert(2, vec![2]).await.unwrap(); - tree.insert(3, vec![3]).await.unwrap(); - tree.insert(4, vec![4]).await.unwrap(); + tree.insert(1, vec![1]).await; + tree.insert(2, vec![2]).await; + tree.insert(3, vec![3]).await; + tree.insert(4, vec![4]).await; let root = tree.root.read().await; match &*root { @@ -816,11 +820,11 @@ mod tests { tree.max_file_size = 100; let large_data = vec![7; 150]; - tree.insert(1, large_data.clone()).await.unwrap(); + tree.insert(1, large_data.clone()).await; let result = tree.get(&1).await.unwrap(); assert_eq!(result, large_data); - tree.insert(2, large_data.clone()).await.unwrap(); + tree.insert(2, large_data.clone()).await; let result = tree.get(&1).await.unwrap(); assert_eq!(result, large_data); @@ -853,53 +857,4 @@ mod tests { ); assert!(loaded_tree.get(&42).await.is_err()); } - - #[tokio::test] - async fn test_save_load_small_tree() { - let tempdir = TempDir::new().unwrap(); - let tree_path = tempdir.path().join("small_tree.bin"); - - let tree = BPlus::::new(2, tempdir.path().into()).unwrap(); - tree.insert(10, vec![1, 2, 3]).await.unwrap(); - tree.insert(20, vec![4, 5, 6]).await.unwrap(); - tree.insert(5, vec![0]).await.unwrap(); - - tree.save(&tree_path).await.unwrap(); - - 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()); - } - - #[tokio::test] - async fn test_save_load_large_tree() { - let tempdir = TempDir::with_prefix("large_load_save").unwrap(); - let tree_path = tempdir.path().join("large_tree.bin"); - let mut tree = BPlus::::new(2, tempdir.path().into()).unwrap(); - tree.max_file_size = 100; - - for i in 0..100000 { - tree.insert(i, vec![(i % 256) as u8; 1]).await.unwrap(); - } - tree.save(&tree_path).await.unwrap(); - - let loaded_tree = BPlus::::load(&tree_path).await.unwrap(); - - use rand::seq::SliceRandom; - use rand::thread_rng; - - let mut rng = thread_rng(); - let mut keys: Vec = (0..10_000).collect(); - keys.shuffle(&mut rng); - - for key in keys.iter().take(100) { - let expected = vec![(*key % 256) as u8; 1]; - assert_eq!(loaded_tree.get(key).await.unwrap(), expected); - } - - assert!(loaded_tree.get(&100_000).await.is_err()); - } } diff --git a/tests/bplus_tests.rs b/tests/bplus_tests.rs index 47e1e62..cfa2a99 100644 --- a/tests/bplus_tests.rs +++ b/tests/bplus_tests.rs @@ -9,7 +9,7 @@ use tempdir::TempDir; async fn test_non_existent_key() { let tempdir = TempDir::new("non_existent").unwrap(); let tree: BPlus = BPlus::new(2, tempdir.path().into()).unwrap(); - tree.insert(1, vec![1]).await.unwrap(); + tree.insert(1, vec![1]).await; assert!(tree.get(&2).await.is_err()); } @@ -18,8 +18,8 @@ async fn test_overwrite_existing_key() { let tempdir = TempDir::new("overwrite").unwrap(); let tree: BPlus = BPlus::new(2, tempdir.path().into()).unwrap(); - tree.insert(1, vec![1]).await.unwrap(); - tree.insert(1, vec![42]).await.unwrap(); + tree.insert(1, vec![1]).await; + tree.insert(1, vec![42]).await; assert_eq!(tree.get(&1).await.unwrap(), vec![42]); } @@ -30,7 +30,7 @@ async fn test_insert_and_find() { let path = PathBuf::new().join(tempdir.path()); let tree: BPlus = BPlus::new(2, path).unwrap(); for i in 1..6 { - tree.insert(i, vec![i as u8; 1]).await.unwrap(); + tree.insert(i, vec![i as u8; 1]).await; } for i in 1..6 { @@ -45,7 +45,7 @@ async fn test_insert_and_find_many_nodes() { let path = PathBuf::new().join(tempdir.path()); let tree: BPlus = BPlus::new(2, path).unwrap(); for i in 1..255 { - tree.insert(i, vec![i as u8; 1]).await.unwrap(); + tree.insert(i, vec![i as u8; 1]).await; } for i in 1..255 { @@ -59,7 +59,7 @@ async fn test_large_data_consecutive_numbers() { let path = PathBuf::new().join(tempdir.path()); let tree: BPlus = BPlus::new(100, path).unwrap(); for i in 1..10000 { - tree.insert(i, vec![i as u8; 1064]).await.unwrap(); + tree.insert(i, vec![i as u8; 1064]).await; } for i in 1..10000 { let a = tree.get(&i).await.unwrap(); @@ -75,7 +75,7 @@ async fn test_large_data() { let mut htable = HashMap::>::new(); for i in 1..10000 { let key = i * 113; - tree.insert(key, vec![key as u8; 1064]).await.unwrap(); + tree.insert(key, vec![key as u8; 1064]).await; htable.insert(key, vec![key as u8; 1064]); } for (key, value) in htable { @@ -88,12 +88,12 @@ async fn test_couple_of_same_keys_inserted() { let tempdir = TempDir::new("8").unwrap(); let tree: BPlus = BPlus::new(2, PathBuf::new().join(tempdir.path())).unwrap(); for i in 1..100 { - tree.insert(i, vec![1u8]).await.unwrap(); + tree.insert(i, vec![1u8]).await; } for i in 1..100 { for j in 1..100 { - tree.insert(i, vec![j as u8]).await.unwrap(); + tree.insert(i, vec![j as u8]).await; } } for i in 1..100 { @@ -114,7 +114,7 @@ async fn test_same_keys_inserted() { } for key in keys.clone() { - tree.insert(key, vec![key as u8]).await.unwrap(); + tree.insert(key, vec![key as u8]).await; } for key in keys { @@ -122,10 +122,10 @@ async fn test_same_keys_inserted() { } let key: usize = rand::random(); - tree.insert(key, vec![0u8]).await.unwrap(); + tree.insert(key, vec![0u8]).await; for i in 1..255 { assert_eq!(vec![i - 1u8], tree.get(&key).await.unwrap()); - tree.insert(key, vec![i]).await.unwrap(); + tree.insert(key, vec![i]).await; } } @@ -135,11 +135,11 @@ async fn test_same_10k_keys_inserted() { let tree: BPlus = BPlus::new(2, PathBuf::new().join(tempdir.path())).unwrap(); for i in 0..10000 { - tree.insert(i, vec![i as u8]).await.unwrap(); + tree.insert(i, vec![i as u8]).await; } for i in 0..10000 { - tree.insert(i, vec![i as u8]).await.unwrap(); + tree.insert(i, vec![i as u8]).await; } for key in 1..10000 { @@ -158,7 +158,7 @@ async fn test_empty_tree() { async fn test_single_entry() { let tempdir = TempDir::new("single").unwrap(); let tree = BPlus::new(2, tempdir.path().into()).unwrap(); - tree.insert(42, vec![1, 2, 3]).await.unwrap(); + tree.insert(42, vec![1, 2, 3]).await; assert_eq!(tree.get(&42).await.unwrap(), vec![1, 2, 3]); } @@ -168,7 +168,7 @@ async fn test_reverse_order_insert() { let tree: BPlus = BPlus::new(3, tempdir.path().into()).unwrap(); for i in (1..100).rev() { - tree.insert(i, vec![i as u8]).await.unwrap(); + tree.insert(i, vec![i as u8]).await; } for i in 1..100 { @@ -182,7 +182,7 @@ async fn test_minimal_degree() { let tree = BPlus::new(1, tempdir.path().into()).unwrap(); for i in 1..=10 { - tree.insert(i, vec![i as u8]).await.unwrap(); + tree.insert(i, vec![i as u8]).await; } assert_eq!(tree.get(&5).await.unwrap(), vec![5]); @@ -194,8 +194,8 @@ async fn test_key_duplication() { let tree = BPlus::new(2, tempdir.path().into()).unwrap(); for _ in 0..10 { - tree.insert(42, vec![1]).await.unwrap(); - tree.insert(42, vec![2]).await.unwrap(); + tree.insert(42, vec![1]).await; + tree.insert(42, vec![2]).await; } assert_eq!(tree.get(&42).await.unwrap(), vec![2]); @@ -206,12 +206,8 @@ async fn test_string_keys() { let tempdir = TempDir::new("string_keys").unwrap(); let tree = BPlus::new(2, tempdir.path().into()).unwrap(); - tree.insert("apple".to_string(), b"fruit".to_vec()) - .await - .unwrap(); - tree.insert("banana".to_string(), b"yellow".to_vec()) - .await - .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()).await.unwrap(), b"fruit"); assert_eq!(tree.get(&"banana".to_string()).await.unwrap(), b"yellow"); @@ -223,7 +219,7 @@ async fn test_stress_1m_entries() { let tree = BPlus::new(100, tempdir.path().into()).unwrap(); for i in 0..1_000_000 { - tree.insert(i, vec![i as u8]).await.unwrap(); + tree.insert(i, vec![i as u8]).await; } for i in 0..1_000_000 { @@ -252,7 +248,7 @@ async fn test_stress_1m_entries_concurrent() { for i in 0..entries_per_task { let key = (task_id * entries_per_task) + i; - tree.insert(key, vec![key as u8]).await.unwrap(); + tree.insert(key, vec![key as u8]).await; } })); } @@ -287,8 +283,48 @@ async fn test_find_nonexistent_after_splits() { let tree = BPlus::new(2, tempdir.path().into()).unwrap(); for i in 0..1000 { - tree.insert(i, vec![1]).await.unwrap(); + tree.insert(i, vec![1]).await; } 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"); + + 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; + + tree.save(&tree_path).await.unwrap(); + + 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()); +} + +#[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(); + + for i in 0..100000 { + tree.insert(i, vec![(i % 256) as u8; 200]).await; + } + tree.save(&tree_path).await.unwrap(); + + let loaded_tree = BPlus::::load(&tree_path).await.unwrap(); + + for key in 0..100000 { + let expected = vec![(key % 256) as u8; 200]; + assert_eq!(loaded_tree.get(&key).await.unwrap(), expected); + } + + assert!(loaded_tree.get(&100_000).await.is_err()); +} From 8298581b6fde9bd3691d8c09a23bfc0ef9d373f0 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Tue, 17 Jun 2025 20:23:47 +0300 Subject: [PATCH 18/22] try fix optimistic_insert --- src/bplus_tree.rs | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index 93e2e0b..336393a 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -580,19 +580,19 @@ impl BPlus { } } - let leaf_lock = match prev_guard.take() { - Some(parent) => { - let pos = last_child_index.unwrap(); - if let Node::Internal(internal) = &*parent { - 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!(); } - None => unreachable!(), }; let mut leaf = leaf_lock.write().await; + drop(prev_guard); if let Node::Leaf(leaf_node) = &mut *leaf { if leaf_node.entries.len() == 2 * self.t - 1 { return Err(()); From a13bb963781a03afcc1ac4b660f760205c20cf89 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Thu, 19 Jun 2025 19:06:29 +0300 Subject: [PATCH 19/22] remove unnecessarycrates --- Cargo.toml | 2 -- 1 file changed, 2 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 0093a2f..f9e9146 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -15,8 +15,6 @@ gnuplot = "0.0.44" approx = "0.5.1" rand = "0.8.5" tokio = { version = "1.44.2", features = ["full"] } -env_logger = "0.11.8" -log = "0.4.27" serde = { version = "1.0", features = ["derive"] } bincode = "1.3" async-recursion = "1.1.1" From 4ee893d6f0cdf9e828b33bddc69e0173a11ae0a3 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Thu, 19 Jun 2025 19:09:37 +0300 Subject: [PATCH 20/22] fix comments --- src/bplus_tree.rs | 16 +++++++++++++--- 1 file changed, 13 insertions(+), 3 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index 336393a..f9528f8 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -25,7 +25,7 @@ const DEFAULT_MAX_FILE_SIZE: u64 = 2 << 20; extern crate chunkfs; -/// Easily serializable version of BPlusTree +/// Serializable version of BPlusTree #[derive(Serialize, Deserialize)] struct SerializableBPlus { t: usize, @@ -55,6 +55,7 @@ struct SerializableLeaf { } impl BPlus { + /// Returns new instance of SerializableBPlus with data from provided BPlus async fn serialize(&self) -> SerializableBPlus { SerializableBPlus { t: self.t, @@ -69,6 +70,7 @@ impl BPlus { impl Node { #[async_recursion] + /// Returns new instance of SerializableNode with data from provided Node async fn serialize(&self) -> SerializableNode { match self { Node::Internal(internal) => { @@ -96,6 +98,7 @@ impl Node { impl Deserialize<'de>> 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))); @@ -228,8 +231,11 @@ pub struct BPlusStorage { 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(); @@ -288,7 +294,9 @@ impl #[allow(dead_code)] 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"); @@ -482,8 +490,6 @@ impl BPlus { for guard in guards { drop(guard); } - - // Ok(()) } #[allow(unused_variables)] @@ -540,9 +546,13 @@ impl BPlus { } /// 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()); From 22fee39e63db4385217b6bea86c406a1d71756b5 Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Thu, 19 Jun 2025 19:23:50 +0300 Subject: [PATCH 21/22] refactor code a little --- src/bplus_tree.rs | 36 ++++++++++++++++-------------------- 1 file changed, 16 insertions(+), 20 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index f9528f8..de7b0df 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -322,14 +322,10 @@ impl BPlus { self.file_number .fetch_add(1, std::sync::atomic::Ordering::SeqCst); self.offset.store(0, std::sync::atomic::Ordering::SeqCst); - *file_guard = File::create( - self.path.join( - self.file_number - .load(std::sync::atomic::Ordering::SeqCst) - .to_string(), - ), - ) - .unwrap(); + 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(); @@ -374,7 +370,7 @@ impl BPlus { // Descent to the leaf loop { let mut current_node = current.write_owned().await; - if let Some(guard) = latch_guard { + if let Some(guard) = latch_guard.take() { drop(guard); latch_guard = None; }; @@ -603,19 +599,19 @@ impl BPlus { let mut leaf = leaf_lock.write().await; drop(prev_guard); - if let Node::Leaf(leaf_node) = &mut *leaf { - 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(()) - } else { + 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(()) } } From 00c871570cc0aba1b478a23e95a56009bb83324e Mon Sep 17 00:00:00 2001 From: kamenkremen Date: Thu, 19 Jun 2025 19:38:46 +0300 Subject: [PATCH 22/22] Add trait alias for BPlusKey and BPlusKeySerializable --- src/bplus_tree.rs | 28 +++++++++++++++------------- 1 file changed, 15 insertions(+), 13 deletions(-) diff --git a/src/bplus_tree.rs b/src/bplus_tree.rs index de7b0df..6bfe1ed 100644 --- a/src/bplus_tree.rs +++ b/src/bplus_tree.rs @@ -23,6 +23,15 @@ 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 @@ -95,9 +104,7 @@ impl Node { } } -impl Deserialize<'de>> - SerializableBPlus -{ +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))); @@ -229,7 +236,7 @@ pub struct BPlusStorage { keys_set: Arc>>, } -impl BPlusStorage { +impl BPlusStorage { /// Creates new instance of B+ tree with given runtime, t and path /// /// runtime is tokio runtime @@ -247,9 +254,7 @@ impl BPlusStorage { } } -impl - Database> for BPlusStorage -{ +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(); @@ -292,7 +297,7 @@ impl } #[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 @@ -615,10 +620,7 @@ impl BPlus { } } -impl< - K: Default + Ord + Clone + Debug + Sized + Serialize + for<'de> Deserialize<'de> + Sync + Send, - > BPlus -{ +impl BPlus { /// Rebuilds links in BPlusTree after loading from file async fn rebuild_links(&self) { let leaves = self.collect_leaves().await; @@ -707,7 +709,7 @@ impl< } } -impl Node { +impl Node { /// Splits node into two and returns new node with it first key fn split(&mut self, t: usize) -> (Link, Arc) { match self {