diff --git a/Cargo.lock b/Cargo.lock index 9833f240dc..741205b934 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2414,6 +2414,7 @@ dependencies = [ "rstest", "serde", "serde_json", + "smallvec", "tempfile", ] diff --git a/Cargo.toml b/Cargo.toml index e7ef64a0e9..7c62df9154 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -205,6 +205,7 @@ serde_cbor = "0.11.2" serde_variant = "0.1.3" sha2 = "0.10.8" serde_json = { version = "~1.0", features = ["preserve_order"] } +smallvec = "1.15.0" strum = { version = "0.26.3", features = ["derive"] } tap = "1.0.1" tar = "0.4.41" diff --git a/lib/collection/Cargo.toml b/lib/collection/Cargo.toml index 7c72cc8a44..92d7e04ddf 100644 --- a/lib/collection/Cargo.toml +++ b/lib/collection/Cargo.toml @@ -43,7 +43,7 @@ hashring = "0.3.6" tinyvec = { version = "1.9.0", features = ["alloc", "latest_stable_rust"] } bitvec = { workspace = true } lazy_static = "1.5.0" -smallvec = "1.15.0" +smallvec = { workspace = true } tokio = { workspace = true } tokio-util = { workspace = true } diff --git a/lib/gridstore/Cargo.toml b/lib/gridstore/Cargo.toml index 437fc25912..dfd99467ff 100644 --- a/lib/gridstore/Cargo.toml +++ b/lib/gridstore/Cargo.toml @@ -17,6 +17,7 @@ ahash = { workspace = true } memmap2 = { workspace = true } serde_json = { workspace = true } serde = { workspace = true } +smallvec = { workspace = true } parking_lot = { workspace = true } tempfile = { workspace = true } lz4_flex = { version = "0.11.3", default-features = false } diff --git a/lib/gridstore/src/gridstore.rs b/lib/gridstore/src/gridstore.rs index 7581fb1870..0563db8874 100644 --- a/lib/gridstore/src/gridstore.rs +++ b/lib/gridstore/src/gridstore.rs @@ -731,43 +731,40 @@ mod tests { fn test_update_single_payload() { let (_dir, mut storage) = empty_storage(); - let mut payload = Payload::default(); - payload.0.insert( - "key".to_string(), - serde_json::Value::String("value".to_string()), - ); let hw_counter = HardwareCounterCell::new(); let hw_counter_ref = hw_counter.ref_payload_io_write_counter(); + let put_payload = + |storage: &mut Gridstore, payload_value: &str, expected_block_offset: u32| { + let mut payload = Payload::default(); + payload.0.insert( + "key".to_string(), + serde_json::Value::String(payload_value.to_string()), + ); - storage.put_value(0, &payload, hw_counter_ref).unwrap(); - assert_eq!(storage.pages.len(), 1); - assert_eq!(storage.tracker.read().mapping_len(), 1); + storage.put_value(0, &payload, hw_counter_ref).unwrap(); + assert_eq!(storage.pages.len(), 1); + assert_eq!(storage.tracker.read().mapping_len(), 1); - let page_mapping = storage.get_pointer(0).unwrap(); - assert_eq!(page_mapping.page_id, 0); // first page - assert_eq!(page_mapping.block_offset, 0); // first cell + let page_mapping = storage.get_pointer(0).unwrap(); + assert_eq!(page_mapping.page_id, 0); // first page + assert_eq!(page_mapping.block_offset, expected_block_offset); - let hw_counter = HardwareCounterCell::new(); - let stored_payload = storage.get_value(0, &hw_counter); - assert!(stored_payload.is_some()); - assert_eq!(stored_payload.unwrap(), payload); + let hw_counter = HardwareCounterCell::new(); + let stored_payload = storage.get_value(0, &hw_counter); + assert!(stored_payload.is_some()); + assert_eq!(stored_payload.unwrap(), payload); + }; - // update payload - let mut updated_payload = Payload::default(); - updated_payload.0.insert( - "key".to_string(), - serde_json::Value::String("updated".to_string()), - ); + put_payload(&mut storage, "value", 0); - storage - .put_value(0, &updated_payload, hw_counter_ref) - .unwrap(); - assert_eq!(storage.pages.len(), 1); - assert_eq!(storage.tracker.read().mapping_len(), 1); + put_payload(&mut storage, "updated", 1); - let stored_payload = storage.get_value(0, &hw_counter); - assert!(stored_payload.is_some()); - assert_eq!(stored_payload.unwrap(), updated_payload); + put_payload(&mut storage, "updated again", 2); + + storage.flush().unwrap(); + + // First block offset should be available again, so we can reuse it + put_payload(&mut storage, "updated after flush", 0); } #[test] diff --git a/lib/gridstore/src/tracker.rs b/lib/gridstore/src/tracker.rs index ab1eb8965f..0d3f7953be 100644 --- a/lib/gridstore/src/tracker.rs +++ b/lib/gridstore/src/tracker.rs @@ -6,6 +6,7 @@ use memory::madvise::{Advice, AdviceSetting}; use memory::mmap_ops::{ create_and_ensure_length, open_write_mmap, transmute_from_u8, transmute_to_u8, }; +use smallvec::SmallVec; pub type PointOffset = u32; pub type BlockOffset = u32; @@ -35,28 +36,57 @@ impl ValuePointer { } } -#[derive(Debug)] -enum PointerUpdate { - Set(ValuePointer), - Unset(ValuePointer), +#[derive(Debug, Default)] +struct PointerUpdates { + /// Whether the latest pointer is set (`true`) or unset (`false`). + /// If this is `true`, then history must have at least one element. + latest_is_set: bool, + /// List of pointers where the value has been written + history: SmallVec<[ValuePointer; 1]>, } -impl PointerUpdate { +impl PointerUpdates { + /// Set the current latest pointer + fn set(&mut self, pointer: ValuePointer) { + self.history.push(pointer); + self.latest_is_set = true; + } + + /// Mark this pointer as pending for freeing + fn unset(&mut self, pointer: ValuePointer) { + // Prevent duplicating pointers to free + if self.history.last() != Some(&pointer) { + self.history.push(pointer); + } + self.latest_is_set = false; + } + #[cfg(test)] fn is_set(&self) -> bool { - match self { - PointerUpdate::Set(_) => true, - PointerUpdate::Unset(_) => false, - } + self.latest_is_set } /// Set is Some, Unset is None - fn to_option(&self) -> Option { - match self { - PointerUpdate::Set(pointer) => Some(*pointer), - PointerUpdate::Unset(_) => None, + fn latest(&self) -> Option { + if self.latest_is_set { + self.history.last().copied() + } else { + None } } + + /// Returns pointers that need to be freed, i.e. They have been written, and are no longer needed + fn into_outdated_pointers(self) -> impl Iterator { + let take = if self.latest_is_set { + // all but the latest one + self.history.len().saturating_sub(1) + } else { + // all of them + self.history.len() + }; + + self.history.into_iter().take(take) + } } #[derive(Debug, Default, Clone)] @@ -75,7 +105,7 @@ pub struct Tracker { /// Updates that haven't been flushed /// /// When flushing, these updates get written into the mmap and flushed at once. - pending_updates: AHashMap, + pending_updates: AHashMap, /// The maximum pointer offset in the tracker (updated in memory). next_pointer_offset: PointOffset, @@ -143,9 +173,9 @@ impl Tracker { // Write pending updates from memory let mut pending_updates = std::mem::take(&mut self.pending_updates); let mut old_pointers = Vec::new(); - for (point_offset, update) in pending_updates.drain() { - match update { - PointerUpdate::Set(new_pointer) => { + for (point_offset, updates) in pending_updates.drain() { + match updates.latest() { + Some(new_pointer) => { if let Some(old_pointer) = self.get_raw(point_offset).and_then(|pointer| *pointer) { @@ -155,12 +185,12 @@ impl Tracker { // write the new pointer self.persist_pointer(point_offset, Some(new_pointer)); } - PointerUpdate::Unset(old_pointer) => { - old_pointers.push(old_pointer); - // write the new pointer + None => { + // write the new None pointer self.persist_pointer(point_offset, None); } } + old_pointers.extend(updates.into_outdated_pointers()); } // increment header count if necessary self.persist_pointer_count(); @@ -262,7 +292,7 @@ impl Tracker { pub fn get(&self, point_offset: PointOffset) -> Option { self.pending_updates .get(&point_offset) - .map(PointerUpdate::to_option) + .map(PointerUpdates::latest) // if the value is not in the pending updates, check the mmap .or_else(|| self.get_raw(point_offset).copied()) .flatten() @@ -280,7 +310,9 @@ impl Tracker { pub fn set(&mut self, point_offset: PointOffset, value_pointer: ValuePointer) { self.pending_updates - .insert(point_offset, PointerUpdate::Set(value_pointer)); + .entry(point_offset) + .or_default() + .set(value_pointer); self.next_pointer_offset = self.next_pointer_offset.max(point_offset + 1); } @@ -290,7 +322,9 @@ impl Tracker { if let Some(pointer) = pointer_opt { self.pending_updates - .insert(point_offset, PointerUpdate::Unset(pointer)); + .entry(point_offset) + .or_default() + .unset(pointer); } pointer_opt diff --git a/lib/segment/Cargo.toml b/lib/segment/Cargo.toml index 6edd8a4836..08ca405ea6 100644 --- a/lib/segment/Cargo.toml +++ b/lib/segment/Cargo.toml @@ -91,7 +91,7 @@ indexmap = { workspace = true } ahash = { workspace = true } self_cell = "1.2.0" sha2 = { workspace = true } -smallvec = "1.15.0" +smallvec = { workspace = true } is_sorted = "0.1.1" strum = { workspace = true } byteorder = { workspace = true }