From 0ea614812818def3dc4143163860d388f7b86508 Mon Sep 17 00:00:00 2001 From: meh Date: Mon, 3 Aug 2026 01:37:56 +0700 Subject: [PATCH] feat: gate synchronized page access --- Cargo.toml | 1 + README.md | 10 ++- src/in_memory/pages.rs | 128 +++++++++++++++++++++++++++++++++++++ src/table/vacuum/vacuum.rs | 18 ++---- 4 files changed, 143 insertions(+), 14 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 93bf735..9f2214f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -17,6 +17,7 @@ categories = ["database-implementations", "data-structures", "caching"] [features] perf_measurements = ["dep:performance_measurement", "dep:performance_measurement_codegen"] s3-support = ["dep:rusty-s3", "dep:url", "dep:reqwest", "dep:walkdir", "worktable_codegen/s3-support"] +safe-concurrent-reads = [] # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html diff --git a/README.md b/README.md index d2d09f3..bc5c78d 100644 --- a/README.md +++ b/README.md @@ -46,6 +46,15 @@ S3 support layers *on top of* the disk engine rather than replacing it. worktable = { version = "0.9", features = ["s3-support"] } # S3 sync, optional ``` +### Concurrent-read safety mode + +For workloads that require reads to be synchronized with in-place row updates +and ghost publication, enable `safe-concurrent-reads`. This opt-in correctness +mode places a table-wide read/write barrier around archived page-byte access. +It prevents readers from observing an in-progress mutation, but serializes page +writes and blocks reads while any page mutation is active. The default is off +to preserve the existing low-latency path. + ## Relationship to `data_bucket` WorkTable is built on [`data_bucket`](https://crates.io/crates/data_bucket), which @@ -396,4 +405,3 @@ enum WorkTableError Check out - [Examples](./examples) - diff --git a/src/in_memory/pages.rs b/src/in_memory/pages.rs index 08a41fb..8d0cbab 100644 --- a/src/in_memory/pages.rs +++ b/src/in_memory/pages.rs @@ -36,6 +36,11 @@ pub struct DataPages where Row: StorableRow, { + /// Serializes access to archived page bytes when `safe-concurrent-reads` + /// is enabled. The default build omits this barrier entirely. + #[cfg(feature = "safe-concurrent-reads")] + page_access: RwLock<()>, + /// Pages vector. Currently, not lock free. pages: RwLock::WrappedRow, DATA_LENGTH>>>>, @@ -68,6 +73,8 @@ where { pub fn new() -> Self { Self { + #[cfg(feature = "safe-concurrent-reads")] + page_access: RwLock::new(()), // We are starting ID's from `1` because `0`'s page in file is info page. pages: RwLock::new(vec![Arc::new(Data::new(1.into()))]), empty_links: EmptyLinkRegistry::::default(), @@ -85,6 +92,8 @@ where } else { let last_page_id = vec.len(); Self { + #[cfg(feature = "safe-concurrent-reads")] + page_access: RwLock::new(()), pages: RwLock::new(vec), empty_links: EmptyLinkRegistry::default(), empty_pages: Default::default(), @@ -104,6 +113,8 @@ where let general_row = ::WrappedRow::from_inner(row); if let Some(link) = self.empty_links.pop_max() { + #[cfg(feature = "safe-concurrent-reads")] + let _page_access = self.page_access.write(); let pages = self.pages.read(); let current_page: usize = page_id_mapper(link.page_id.into()); let page = &pages[current_page]; @@ -129,6 +140,8 @@ where loop { let (link, tried_page) = { + #[cfg(feature = "safe-concurrent-reads")] + let _page_access = self.page_access.write(); let pages = self.pages.read(); let current_page = page_id_mapper(self.current_page_id.load(Ordering::Acquire) as usize); let page = &pages[current_page]; @@ -197,6 +210,8 @@ where }; if let Some(page_id) = page_id { + #[cfg(feature = "safe-concurrent-reads")] + let _page_access = self.page_access.write(); let pages = self.pages.read(); let index = page_id_mapper(page_id.into()); let page = pages[index].clone(); @@ -221,6 +236,8 @@ where Portable + Deserialize<::WrappedRow, HighDeserializer>, { let link = link.into(); + #[cfg(feature = "safe-concurrent-reads")] + let _page_access = self.page_access.read(); let pages = self.pages.read(); let page = pages .get(page_id_mapper(link.page_id.into())) @@ -235,6 +252,8 @@ where <::WrappedRow as Archive>::Archived: Portable + Deserialize<::WrappedRow, HighDeserializer>, { + #[cfg(feature = "safe-concurrent-reads")] + let _page_access = self.page_access.read(); let pages = self.pages.read(); let page = pages .get(page_id_mapper(link.page_id.into())) @@ -252,6 +271,8 @@ where <::WrappedRow as Archive>::Archived: Portable + Deserialize<::WrappedRow, HighDeserializer>, { + #[cfg(feature = "safe-concurrent-reads")] + let _page_access = self.page_access.read(); let pages = self.pages.read(); let page = pages .get(page_id_mapper(link.page_id.into())) @@ -272,6 +293,8 @@ where Row: Archive + for<'a> Serialize, Share>, rkyv::rancor::Error>>, Op: Fn(&<::WrappedRow as Archive>::Archived) -> Res, { + #[cfg(feature = "safe-concurrent-reads")] + let _page_access = self.page_access.read(); let pages = self.pages.read(); let page = pages .get::(page_id_mapper(link.page_id.into())) @@ -289,6 +312,8 @@ where <::WrappedRow as Archive>::Archived: Portable, Op: FnMut(&mut <::WrappedRow as Archive>::Archived) -> Res, { + #[cfg(feature = "safe-concurrent-reads")] + let _page_access = self.page_access.write(); let pages = self.pages.read(); let page = pages .get(page_id_mapper(link.page_id.into())) @@ -314,6 +339,8 @@ where ::WrappedRow: Archive + for<'a> Serialize, Share>, rkyv::rancor::Error>>, { + #[cfg(feature = "safe-concurrent-reads")] + let _page_access = self.page_access.write(); let pages = self.pages.read(); let page = pages .get(page_id_mapper(link.page_id.into())) @@ -339,6 +366,8 @@ where } pub fn select_raw(&self, link: Link) -> Result, ExecutionError> { + #[cfg(feature = "safe-concurrent-reads")] + let _page_access = self.page_access.read(); let pages = self.pages.read(); let page = pages .get(page_id_mapper(link.page_id.into())) @@ -395,6 +424,8 @@ where } pub fn get_bytes(&self) -> Vec<([u8; DATA_LENGTH], u32)> { + #[cfg(feature = "safe-concurrent-reads")] + let _page_access = self.page_access.read(); let pages = self.pages.read(); pages .iter() @@ -402,6 +433,52 @@ where .collect() } + pub(crate) fn reset_page(&self, page_id: PageId) -> Result<(), ExecutionError> { + #[cfg(feature = "safe-concurrent-reads")] + let _page_access = self.page_access.write(); + let pages = self.pages.read(); + let page = pages + .get(page_id_mapper(page_id.into())) + .ok_or(ExecutionError::PageNotFound(page_id))?; + page.reset(); + Ok(()) + } + + /// Copies a row to another page and retires the old copy as one page-byte + /// mutation. In strict mode, readers cannot observe either write midway. + pub(crate) unsafe fn move_row_for_vacuum( + &self, + from_link: Link, + to_page_id: PageId, + ) -> Result<(Vec, Link), ExecutionError> + where + <::WrappedRow as Archive>::Archived: ArchivedRowWrapper + Portable, + { + #[cfg(feature = "safe-concurrent-reads")] + let _page_access = self.page_access.write(); + let pages = self.pages.read(); + let from_page = pages + .get(page_id_mapper(from_link.page_id.into())) + .ok_or(ExecutionError::PageNotFound(from_link.page_id))?; + let to_page = pages + .get(page_id_mapper(to_page_id.into())) + .ok_or(ExecutionError::PageNotFound(to_page_id))?; + + let raw_data = from_page + .get_raw_row(from_link) + .map_err(ExecutionError::DataPageError)?; + let archived = unsafe { + from_page + .get_mut_row_ref(from_link) + .map_err(ExecutionError::DataPageError)? + .unseal_unchecked() + }; + archived.set_in_vacuum_process(); + let new_link = to_page.save_raw_row(&raw_data).map_err(ExecutionError::DataPageError)?; + + Ok((raw_data, new_link)) + } + pub fn get_page_count(&self) -> usize { self.pages.read().len() } @@ -466,7 +543,11 @@ mod tests { use std::collections::HashSet; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; + #[cfg(feature = "safe-concurrent-reads")] + use std::sync::mpsc; use std::thread; + #[cfg(feature = "safe-concurrent-reads")] + use std::time::Duration; use std::time::Instant; use parking_lot::RwLock; @@ -688,6 +769,53 @@ mod tests { ); } + #[cfg(feature = "safe-concurrent-reads")] + #[test] + fn safe_concurrent_read_waits_for_complete_ghost_publication() { + let pages = Arc::new(DataPages::::new()); + let link = pages.insert(TestRow { a: 0, b: 0 }).unwrap(); + let (mutation_started_tx, mutation_started_rx) = mpsc::channel(); + let (finish_mutation_tx, finish_mutation_rx) = mpsc::channel(); + + let writer_pages = pages.clone(); + let writer = thread::spawn(move || unsafe { + writer_pages + .with_mut_ref(link, |archived| { + archived.inner.a = 1.into(); + mutation_started_tx.send(()).unwrap(); + finish_mutation_rx.recv().unwrap(); + archived.inner.b = 1.into(); + archived.unghost(); + }) + .unwrap(); + }); + + mutation_started_rx.recv().unwrap(); + + let reader_pages = pages.clone(); + let (reader_started_tx, reader_started_rx) = mpsc::channel(); + let (read_result_tx, read_result_rx) = mpsc::channel(); + let reader = thread::spawn(move || { + reader_started_tx.send(()).unwrap(); + read_result_tx.send(reader_pages.select_non_ghosted(link)).unwrap(); + }); + + reader_started_rx.recv().unwrap(); + assert!( + read_result_rx.recv_timeout(Duration::from_millis(50)).is_err(), + "reader passed the page barrier while publication was incomplete" + ); + + finish_mutation_tx.send(()).unwrap(); + assert_eq!( + read_result_rx.recv_timeout(Duration::from_secs(1)).unwrap(), + Ok(TestRow { a: 1, b: 1 }) + ); + + writer.join().unwrap(); + reader.join().unwrap(); + } + #[test] fn update() { let pages = DataPages::::new(); diff --git a/src/table/vacuum/vacuum.rs b/src/table/vacuum/vacuum.rs index 1cb3d8b..7de5805 100644 --- a/src/table/vacuum/vacuum.rs +++ b/src/table/vacuum/vacuum.rs @@ -207,13 +207,11 @@ where } fn free_page(&self, page_id: PageId) { - let p = self.data_pages.get_page(page_id).expect("should exist as called"); - p.reset() + self.data_pages.reset_page(page_id).expect("should exist as called") } async fn move_data_from(&self, from: PageId, to: PageId) -> (bool, bool) { let to_page = self.data_pages.get_page(to).expect("should exist as link exists"); - let from_page = self.data_pages.get_page(from).expect("should exist as link exists"); let to_free_space = to_page.free_space(); let page_start = OffsetEqLink::<_>(Link { @@ -268,17 +266,11 @@ where self.lock_manager.remove_with_lock_check(&pk); continue; } - let raw_data = from_page - .get_raw_row(from_link.0) - .expect("link is not bigger than free offset"); - unsafe { + let (raw_data, new_link) = unsafe { self.data_pages - .with_mut_ref(from_link.0, |r| r.set_in_vacuum_process()) - .expect("link should be valid") - } - let new_link = to_page - .save_raw_row(&raw_data) - .expect("page is not full as checked on links collection"); + .move_row_for_vacuum(from_link.0, to) + .expect("links and destination capacity were checked") + }; self.update_index_after_move(pk.clone(), from_link.0, new_link, raw_data); lock.unlock();