From 924d2b652ed010718f51a15f0e6f0b7e5f00fb05 Mon Sep 17 00:00:00 2001 From: meh Date: Mon, 3 Aug 2026 01:52:19 +0700 Subject: [PATCH] fix: replace vacuum delay with epoch grace period --- benches/cases/simple.rs | 23 ++++++ src/in_memory/empty_link_registry.rs | 103 ++++++++++++++++++++++++--- src/in_memory/pages.rs | 78 ++++++++++++++++++-- src/table/mod.rs | 22 ++++-- src/table/vacuum/vacuum.rs | 7 +- 5 files changed, 210 insertions(+), 23 deletions(-) diff --git a/benches/cases/simple.rs b/benches/cases/simple.rs index 4ce590c..97d1798 100644 --- a/benches/cases/simple.rs +++ b/benches/cases/simple.rs @@ -90,6 +90,28 @@ fn delete(c: &mut Criterion) { }); } +fn delete_insert_reuse(c: &mut Criterion) { + let rt = Runtime::new().unwrap(); + let table = Arc::new(SimpleWorkTable::default()); + let initial = SimpleRow { + id: table.get_next_pk().into(), + value: fastrand::u64(..), + }; + let initial_pk = table.insert(initial).unwrap(); + rt.block_on(table.delete(initial_pk)).unwrap(); + + c.bench_function("simple_delete_insert_reuse", |b| { + b.to_async(&rt).iter(|| async { + let row = SimpleRow { + id: table.get_next_pk().into(), + value: fastrand::u64(..), + }; + let pk = table.insert(black_box(row)).unwrap(); + black_box(table.delete(pk).await) + }) + }); +} + fn upsert_insert(c: &mut Criterion) { let rt = Runtime::new().unwrap(); let table = Arc::new(SimpleWorkTable::default()); @@ -195,6 +217,7 @@ criterion_group! { select_by_pk, update, delete, + delete_insert_reuse, upsert_insert, upsert_update, batch_insert, diff --git a/src/in_memory/empty_link_registry.rs b/src/in_memory/empty_link_registry.rs index 3030420..880b6b1 100644 --- a/src/in_memory/empty_link_registry.rs +++ b/src/in_memory/empty_link_registry.rs @@ -1,4 +1,5 @@ -use std::sync::atomic::{AtomicU32, Ordering}; +use std::ops::Deref; +use std::sync::atomic::{AtomicU32, AtomicU64, AtomicUsize, Ordering}; use data_bucket::Link; use data_bucket::page::PageId; @@ -85,6 +86,38 @@ pub struct EmptyLinkRegistry { pub(crate) op_lock: FairMutex<()>, vacuum_lock: tokio::sync::Mutex<()>, + reuse_epoch: AtomicU64, + active_reusers: [AtomicUsize; 2], + reuse_quiescent: tokio::sync::Notify, +} + +/// A free link reserved by an insert operation. +/// +/// Dropping this lease announces that the link is no longer in flight. Vacuum +/// advances the reuse epoch and waits for all leases from the preceding epoch +/// before it starts moving or resetting pages. +#[derive(Debug)] +#[must_use = "the reservation must live until the free-link write finishes"] +pub(crate) struct EmptyLinkReservation<'a, const DATA_LENGTH: usize = DATA_INNER_LENGTH> { + registry: &'a EmptyLinkRegistry, + link: Link, + epoch: usize, +} + +impl Deref for EmptyLinkReservation<'_, DATA_LENGTH> { + type Target = Link; + + fn deref(&self) -> &Self::Target { + &self.link + } +} + +impl Drop for EmptyLinkReservation<'_, DATA_LENGTH> { + fn drop(&mut self) { + if self.registry.active_reusers[self.epoch].fetch_sub(1, Ordering::AcqRel) == 1 { + self.registry.reuse_quiescent.notify_one(); + } + } } impl Default for EmptyLinkRegistry { @@ -96,6 +129,9 @@ impl Default for EmptyLinkRegistry { sum_links_len: Default::default(), op_lock: Default::default(), vacuum_lock: Default::default(), + reuse_epoch: AtomicU64::new(0), + active_reusers: [AtomicUsize::new(0), AtomicUsize::new(0)], + reuse_quiescent: tokio::sync::Notify::new(), } } } @@ -164,20 +200,36 @@ impl EmptyLinkRegistry { self.insert_link(index_ord_link); } + /// Removes and returns the largest free link. pub fn pop_max(&self) -> Option { - if self.vacuum_lock.try_lock().is_err() { - return None; - } + self.reserve_max().map(|reservation| *reservation) + } + + /// Reserves the largest free link for internal page reuse. + /// + /// The returned lease must remain alive until the caller has published the + /// row written into the link. Vacuum admission is closed while the lease is + /// registered, so a concurrent epoch advance cannot miss it. + pub(crate) fn reserve_max(&self) -> Option> { + let _vacuum_admission = self.vacuum_lock.try_lock().ok()?; let _g = self.op_lock.lock(); let mut iter = self.length_ord_links.iter().rev(); let (_, max_length_link) = iter.next()?; + let max_length_link = *max_length_link; drop(iter); - self.remove_link(*max_length_link); + self.remove_link(max_length_link); - Some(*max_length_link) + let epoch = self.reuse_epoch.load(Ordering::Acquire) as usize & 1; + self.active_reusers[epoch].fetch_add(1, Ordering::AcqRel); + + Some(EmptyLinkReservation { + registry: self, + link: max_length_link, + epoch, + }) } pub fn iter(&self) -> impl Iterator + '_ { @@ -188,8 +240,17 @@ impl EmptyLinkRegistry { self.sum_links_len.load(Ordering::Acquire) } + /// Closes free-link reuse admission, advances the epoch, and waits for all + /// reservations admitted in the previous epoch to finish. pub async fn lock_vacuum(&self) -> tokio::sync::MutexGuard<'_, ()> { - self.vacuum_lock.lock().await + let guard = self.vacuum_lock.lock().await; + let previous_epoch = self.reuse_epoch.fetch_add(1, Ordering::AcqRel) as usize & 1; + + while self.active_reusers[previous_epoch].load(Ordering::Acquire) != 0 { + self.reuse_quiescent.notified().await; + } + + guard } } @@ -422,7 +483,7 @@ mod tests { fn test_empty_registry() { let registry = EmptyLinkRegistry::::default(); - assert_eq!(registry.pop_max(), None); + assert!(registry.pop_max().is_none()); assert_eq!(registry.iter().count(), 0); } @@ -489,4 +550,30 @@ mod tests { ); assert_eq!(popped_after_unlock.unwrap().length, 100); } + + #[tokio::test] + async fn test_lock_vacuum_waits_for_previous_epoch_reuser() { + let registry = EmptyLinkRegistry::::default(); + registry.push(Link { + page_id: 1.into(), + offset: 0, + length: 100, + }); + + let reservation = registry.reserve_max().unwrap(); + let vacuum = registry.lock_vacuum(); + tokio::pin!(vacuum); + + assert!( + tokio::time::timeout(std::time::Duration::from_millis(50), &mut vacuum) + .await + .is_err(), + "vacuum passed the grace period while a prior-epoch link was in flight" + ); + + drop(reservation); + let _guard = tokio::time::timeout(std::time::Duration::from_secs(1), vacuum) + .await + .expect("vacuum should resume when the prior epoch becomes quiescent"); + } } diff --git a/src/in_memory/pages.rs b/src/in_memory/pages.rs index 08a41fb..9728693 100644 --- a/src/in_memory/pages.rs +++ b/src/in_memory/pages.rs @@ -17,7 +17,7 @@ use std::{ sync::atomic::{AtomicU32, AtomicU64, Ordering}, }; -use crate::in_memory::empty_link_registry::EmptyLinkRegistry; +use crate::in_memory::empty_link_registry::{EmptyLinkRegistry, EmptyLinkReservation}; use crate::prelude::ArchivedRowWrapper; use crate::{ in_memory::{ @@ -31,6 +31,17 @@ fn page_id_mapper(page_id: usize) -> usize { page_id - 1usize } +pub(crate) struct InsertedRow<'a, const DATA_LENGTH: usize> { + link: Link, + _reuse_reservation: Option>, +} + +impl InsertedRow<'_, DATA_LENGTH> { + pub(crate) fn link(&self) -> Link { + self.link + } +} + #[derive(Debug)] pub struct DataPages where @@ -96,6 +107,15 @@ where } pub fn insert(&self, row: Row) -> Result + where + Row: Archive + for<'a> Serialize, Share>, rkyv::rancor::Error>>, + ::WrappedRow: + Archive + for<'a> Serialize, Share>, rkyv::rancor::Error>>, + { + Ok(self.insert_with_reservation(row)?.link()) + } + + pub(crate) fn insert_with_reservation(&self, row: Row) -> Result, ExecutionError> where Row: Archive + for<'a> Serialize, Share>, rkyv::rancor::Error>>, ::WrappedRow: @@ -103,7 +123,8 @@ where { let general_row = ::WrappedRow::from_inner(row); - if let Some(link) = self.empty_links.pop_max() { + if let Some(reservation) = self.empty_links.reserve_max() { + let link = *reservation; let pages = self.pages.read(); let current_page: usize = page_id_mapper(link.page_id.into()); let page = &pages[current_page]; @@ -113,7 +134,10 @@ where if let Some(l) = left_link { self.empty_links.push(l); } - return Ok(link); + return Ok(InsertedRow { + link, + _reuse_reservation: Some(reservation), + }); } Err(e) => match e { DataExecutionError::InvalidLink => { @@ -138,7 +162,10 @@ where match link { Ok(link) => { self.row_count.fetch_add(1, Ordering::Relaxed); - return Ok(link); + return Ok(InsertedRow { + link, + _reuse_reservation: None, + }); } Err(e) => match e { DataExecutionError::PageIsFull { .. } => { @@ -170,12 +197,27 @@ where ::WrappedRow: Archive + for<'a> Serialize, Share>, rkyv::rancor::Error>>, { - let link = self.insert(row.clone())?; + let (inserted, bytes) = self.insert_cdc_with_reservation(row)?; + Ok((inserted.link(), bytes)) + } + + pub(crate) fn insert_cdc_with_reservation( + &self, + row: Row, + ) -> Result<(InsertedRow<'_, DATA_LENGTH>, Vec), ExecutionError> + where + Row: Archive + + for<'a> Serialize, Share>, rkyv::rancor::Error>> + + Clone, + ::WrappedRow: + Archive + for<'a> Serialize, Share>, rkyv::rancor::Error>>, + { + let inserted = self.insert_with_reservation(row.clone())?; let general_row = ::WrappedRow::from_inner(row); let bytes = rkyv::to_bytes(&general_row) .expect("should be ok as insert not failed") .into_vec(); - Ok((link, bytes)) + Ok((inserted, bytes)) } fn add_next_page(&self, tried_page: usize) { @@ -728,6 +770,30 @@ mod tests { assert_eq!(pages.select(link).unwrap(), TestRow { a: 10, b: 20 }) } + #[tokio::test] + async fn reused_link_reservation_lives_through_caller_publication() { + let pages = DataPages::::new(); + let old_link = pages.insert(TestRow { a: 10, b: 20 }).unwrap(); + pages.delete(old_link).unwrap(); + + let inserted = pages.insert_with_reservation(TestRow { a: 30, b: 40 }).unwrap(); + assert_eq!(inserted.link(), old_link); + + let vacuum = pages.empty_links.lock_vacuum(); + tokio::pin!(vacuum); + assert!( + tokio::time::timeout(std::time::Duration::from_millis(50), &mut vacuum) + .await + .is_err(), + "vacuum passed the grace period before the caller published the reused link" + ); + + drop(inserted); + let _guard = tokio::time::timeout(std::time::Duration::from_secs(1), vacuum) + .await + .expect("vacuum should resume after publication releases the reused link"); + } + //#[test] fn _bench() { let pages = Arc::new(DataPages::::new()); diff --git a/src/table/mod.rs b/src/table/mod.rs index cb54aaa..248beae 100644 --- a/src/table/mod.rs +++ b/src/table/mod.rs @@ -176,7 +176,11 @@ where LockType: 'static, { let pk = row.get_primary_key().clone(); - let link = self.data.insert(row.clone()).map_err(WorkTableError::PagesError)?; + let inserted = self + .data + .insert_with_reservation(row.clone()) + .map_err(WorkTableError::PagesError)?; + let link = inserted.link(); if self.primary_index.insert_checked(pk.clone(), link).is_none() { self.data.delete(link).map_err(WorkTableError::PagesError)?; return Err(WorkTableError::PrimaryAlreadyExists); @@ -198,6 +202,7 @@ where .with_mut_ref(link, |r| r.unghost()) .map_err(WorkTableError::PagesError)? } + drop(inserted); Ok(pk) } @@ -227,10 +232,11 @@ where { let pk = row.get_primary_key().clone(); - let (link, _) = match self.data.insert_cdc(row.clone()) { + let (inserted, _) = match self.data.insert_cdc_with_reservation(row.clone()) { Ok(result) => result, Err(e) => return (None, Err(WorkTableError::PagesError(e))), }; + let link = inserted.link(); let primary_key_events = self.primary_index.insert_checked_cdc(pk.clone(), link); let Some(primary_key_events) = primary_key_events else { @@ -304,6 +310,7 @@ where return (Some(ack_op), Err(WorkTableError::PagesError(e))); } }; + drop(inserted); let op = Operation::Insert(InsertOperation { id: OperationId::Single(Uuid::now_v7()), @@ -351,7 +358,11 @@ where .get(&pk) .map(|v| v.get().value.into()) .ok_or(WorkTableError::NotFound)?; - let new_link = self.data.insert(row_new.clone()).map_err(WorkTableError::PagesError)?; + let inserted = self + .data + .insert_with_reservation(row_new.clone()) + .map_err(WorkTableError::PagesError)?; + let new_link = inserted.link(); unsafe { self.data .with_mut_ref(new_link, |r| r.unghost()) @@ -373,6 +384,7 @@ where }; } self.data.delete(old_link).map_err(WorkTableError::PagesError)?; + drop(inserted); Ok(pk) } @@ -411,10 +423,11 @@ where }; // Insert new data - if this fails, no events to acknowledge - let (new_link, _) = match self.data.insert_cdc(row_new.clone()) { + let (inserted, _) = match self.data.insert_cdc_with_reservation(row_new.clone()) { Ok(result) => result, Err(e) => return (None, Err(WorkTableError::PagesError(e))), }; + let new_link = inserted.link(); // Unghost the new data - if this fails, we have no events yet to acknowledge unsafe { @@ -497,6 +510,7 @@ where return (Some(ack_op), Err(WorkTableError::PagesError(e))); } }; + drop(inserted); let op = Operation::Insert(InsertOperation { id: OperationId::Single(Uuid::now_v7()), diff --git a/src/table/vacuum/vacuum.rs b/src/table/vacuum/vacuum.rs index 1cb3d8b..7607e8a 100644 --- a/src/table/vacuum/vacuum.rs +++ b/src/table/vacuum/vacuum.rs @@ -2,7 +2,7 @@ use std::collections::VecDeque; use std::fmt::Debug; use std::marker::PhantomData; use std::sync::Arc; -use std::time::{Duration, Instant}; +use std::time::Instant; use data_bucket::Link; use data_bucket::page::PageId; @@ -133,11 +133,8 @@ where let now = Instant::now(); let registry = self.data_pages.empty_links_registry(); - let mut per_page_info = registry.get_per_page_info(); let _registry_lock = registry.lock_vacuum().await; - - // to avoid some rewrites of ops that used link from empty links registry - tokio::time::sleep(Duration::from_millis(100)).await; + let mut per_page_info = registry.get_per_page_info(); per_page_info.sort_by_key(|l| OrderedFloat(l.filled_empty_ratio)); let initial_bytes_freed: u64 = per_page_info.iter().map(|i| i.empty_bytes as u64).sum();