Skip to content
This repository was archived by the owner on Aug 3, 2026. It is now read-only.
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
10 changes: 9 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -396,4 +405,3 @@ enum WorkTableError

Check out - [Examples](./examples)


128 changes: 128 additions & 0 deletions src/in_memory/pages.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,11 @@ pub struct DataPages<Row, const DATA_LENGTH: usize = DATA_INNER_LENGTH>
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<Vec<Arc<Data<<Row as StorableRow>::WrappedRow, DATA_LENGTH>>>>,

Expand Down Expand Up @@ -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::<DATA_LENGTH>::default(),
Expand All @@ -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(),
Expand All @@ -104,6 +113,8 @@ where
let general_row = <Row as StorableRow>::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];
Expand All @@ -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];
Expand Down Expand Up @@ -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();
Expand All @@ -221,6 +236,8 @@ where
Portable + Deserialize<<Row as StorableRow>::WrappedRow, HighDeserializer<rkyv::rancor::Error>>,
{
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()))
Expand All @@ -235,6 +252,8 @@ where
<<Row as StorableRow>::WrappedRow as Archive>::Archived:
Portable + Deserialize<<Row as StorableRow>::WrappedRow, HighDeserializer<rkyv::rancor::Error>>,
{
#[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()))
Expand All @@ -252,6 +271,8 @@ where
<<Row as StorableRow>::WrappedRow as Archive>::Archived:
Portable + Deserialize<<Row as StorableRow>::WrappedRow, HighDeserializer<rkyv::rancor::Error>>,
{
#[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()))
Expand All @@ -272,6 +293,8 @@ where
Row: Archive + for<'a> Serialize<Strategy<Serializer<AlignedVec, ArenaHandle<'a>, Share>, rkyv::rancor::Error>>,
Op: Fn(&<<Row as StorableRow>::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::<usize>(page_id_mapper(link.page_id.into()))
Expand All @@ -289,6 +312,8 @@ where
<<Row as StorableRow>::WrappedRow as Archive>::Archived: Portable,
Op: FnMut(&mut <<Row as StorableRow>::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()))
Expand All @@ -314,6 +339,8 @@ where
<Row as StorableRow>::WrappedRow:
Archive + for<'a> Serialize<Strategy<Serializer<AlignedVec, ArenaHandle<'a>, 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()))
Expand All @@ -339,6 +366,8 @@ where
}

pub fn select_raw(&self, link: Link) -> Result<Vec<u8>, 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()))
Expand Down Expand Up @@ -395,13 +424,61 @@ 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()
.map(|p| (p.get_bytes(), p.free_offset.load(Ordering::Relaxed)))
.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<u8>, Link), ExecutionError>
where
<<Row as StorableRow>::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()
}
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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::<TestRow>::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::<TestRow>::new();
Expand Down
18 changes: 5 additions & 13 deletions src/table/vacuum/vacuum.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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();
Expand Down
Loading