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"]
versioned-row-publication = ["worktable_codegen/versioned-row-publication"]

# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html

Expand Down
25 changes: 24 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,30 @@ 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/write publication

The default build preserves the existing lowest-latency page path and requires
applications to exclude reads that overlap page-byte mutation. Applications
that need generated reads to overlap updates, inserts, deletes, and vacuum can
opt into immutable row-version publication:

```toml
[dependencies]
worktable = { version = "0.9", features = ["versioned-row-publication"] }
```

In this mode, generated reads acquire an immutable owned row version instead
of borrowing the mutable archived page image. Writers replace a per-row version
only after a complete page mutation, insert visibility is an atomic lifecycle
transition after every index is installed, and deleted or relocated links are
not reused until readers that could have captured them have drained. Page bytes
remain the persistence image and are internally serialized; range queries are
still non-snapshot reads. The mode intentionally trades memory, an atomic
read-side grace-period counter, and publication bookkeeping for this stronger
concurrent-read contract. See
[`docs/versioned-row-publication.md`](docs/versioned-row-publication.md) for the
protocol and its scope.

## Relationship to `data_bucket`

WorkTable is built on [`data_bucket`](https://crates.io/crates/data_bucket), which
Expand Down Expand Up @@ -396,4 +420,3 @@ enum WorkTableError

Check out - [Examples](./examples)


22 changes: 22 additions & 0 deletions benches/cases/unique_index.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
use criterion::{BatchSize, BenchmarkId, Criterion, Throughput, black_box, criterion_group};
use std::sync::Arc;
use tokio::runtime::Runtime;
use worktable::prelude::SelectQueryExecutor;

use crate::common::*;

Expand Down Expand Up @@ -64,6 +65,26 @@ fn select_by_unique_index(c: &mut Criterion) {
});
}

fn select_by_unique_index_range(c: &mut Criterion) {
let table = UniqueIndexWorkTable::default();

for i in 1..=1000i64 {
let row = UniqueIndexRow {
id: table.get_next_pk().into(),
test: i,
another: i as u64,
};
table.insert(row).unwrap();
}

c.bench_function("unique_index_select_by_test_range", |b| {
b.iter(|| {
let test = fastrand::i64(1..=1000);
black_box(table.select_by_test_range(test..=test).execute())
})
});
}

fn update(c: &mut Criterion) {
let rt = Runtime::new().unwrap();
let table = Arc::new(UniqueIndexWorkTable::default());
Expand Down Expand Up @@ -224,6 +245,7 @@ criterion_group! {
insert,
select_by_pk,
select_by_unique_index,
select_by_unique_index_range,
update,
delete,
upsert_insert,
Expand Down
1 change: 1 addition & 0 deletions codegen/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ repository = "https://github.com/pathscale/WorkTable"

[features]
s3-support = []
versioned-row-publication = []

[lib]
name = "worktable_codegen"
Expand Down
6 changes: 5 additions & 1 deletion codegen/src/generators/in_memory/queries/select.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,9 +29,13 @@ impl InMemoryGenerator {
#column_range_type,
#row_fields_ident>
{
let read_guard = self.0.data.read_guard();
let iter = self.0.primary_index.pk_map
.iter()
.filter_map(|(_, link)| self.0.data.select_non_ghosted(link.0).ok());
.filter_map(move |(_, link)| {
let _read_guard = &read_guard;
self.0.data.select_non_ghosted(link.0).ok()
});

SelectQueryBuilder::new(iter)
}
Expand Down
7 changes: 6 additions & 1 deletion codegen/src/generators/in_memory/table/impls.rs
Original file line number Diff line number Diff line change
Expand Up @@ -110,9 +110,13 @@ impl InMemoryGenerator {
range.start_bound().map(|v| #primary_key_type::from(v.clone())),
range.end_bound().map(|v| #primary_key_type::from(v.clone())),
);
let read_guard = self.0.data.read_guard();
let rows = self.0.primary_index.pk_map
.range(converted_range)
.filter_map(|(_, link)| self.0.data.select_non_ghosted(link.0).ok());
.filter_map(move |(_, link)| {
let _read_guard = &read_guard;
self.0.data.select_non_ghosted(link.0).ok()
});

#pk_sorted_by
}
Expand Down Expand Up @@ -262,6 +266,7 @@ impl InMemoryGenerator {

fn gen_table_iter_inner(&self, func: TokenStream) -> TokenStream {
quote! {
let _read_guard = self.0.data.read_guard();
let first = self.0.primary_index.pk_map.iter().next().map(|(k, v)| (k.clone(), v.0));
let Some((mut k, link)) = first else {
return Ok(())
Expand Down
79 changes: 73 additions & 6 deletions codegen/src/generators/in_memory/table/index_fns.rs
Original file line number Diff line number Diff line change
Expand Up @@ -83,13 +83,39 @@ impl InMemoryGenerator {
row.#row_field_ident.eq(&by)
}
};

Ok(quote! {
pub fn #fn_name(&self, by: #type_) -> Option<#row_ident> {
let select = if cfg!(feature = "versioned-row-publication") {
quote! {
loop {
let link: Link = self.0.indexes.#field_ident
.get(#by)
.map(|kv| kv.get().value.into())?;
if let Ok(row) = self.0.data.select_non_ghosted(link) {
if #predicate_matches {
return Some(row);
}
}

let current_link: Option<Link> = self.0.indexes.#field_ident
.get(#by)
.map(|kv| kv.get().value.into());
if current_link == Some(link) {
return None;
}
}
}
} else {
quote! {
let link: Link = self.0.indexes.#field_ident.get(#by).map(|kv| kv.get().value.into())?;
let row = self.0.data.select_non_ghosted(link).ok()?;
#predicate_matches.then_some(row)
}
};

Ok(quote! {
pub fn #fn_name(&self, by: #type_) -> Option<#row_ident> {
let _read_guard = self.0.data.read_guard();
#select
}
})
}

Expand Down Expand Up @@ -121,10 +147,14 @@ impl InMemoryGenerator {
#column_range_type,
#row_fields_ident>
{
let read_guard = self.0.data.read_guard();
let rows = self.0.indexes.#field_ident
.get(#by)
.into_iter()
.filter_map(|(_, link)| self.0.data.select_non_ghosted(link.0).ok())
.filter_map(move |(_, link)| {
let _read_guard = &read_guard;
self.0.data.select_non_ghosted(link.0).ok()
})
.filter(move |r| &r.#row_field_ident == &by);

SelectQueryBuilder::new(rows)
Expand All @@ -143,21 +173,52 @@ impl InMemoryGenerator {
let type_ = columns_map.get(i).ok_or(syn::Error::new(i.span(), "Row not found"))?;
let fn_name = Ident::new(format!("select_by_{i}_range").as_str(), Span::mixed_site());
let field_ident = &idx.name;
let row_field_ident = &idx.field;
let column_pascal = Ident::new(&i.to_string().to_case(Case::Pascal), Span::mixed_site());

let revalidate = cfg!(feature = "versioned-row-publication");
let (range_bounds, range_arg) = if is_float(type_.to_string().as_str()) {
(
quote! { std::ops::RangeBounds<#type_> },
quote! {
if revalidate {
quote! {
(
predicate_range.0.as_ref().map(|v| OrderedFloat(*v)),
predicate_range.1.as_ref().map(|v| OrderedFloat(*v)),
)
}
} else {
quote! {
(
range.start_bound().map(|v| OrderedFloat(*v)),
range.end_bound().map(|v| OrderedFloat(*v)),
)
}
},
)
} else if revalidate {
(
quote! { std::ops::RangeBounds<#type_> },
quote! { predicate_range.clone() },
)
} else {
(quote! { std::ops::RangeBounds<#type_> }, quote! { range })
};
let predicate_setup = revalidate.then(|| {
quote! {
let predicate_range = (
range.start_bound().cloned(),
range.end_bound().cloned(),
);
}
});
let predicate_filter = revalidate.then(|| {
quote! {
.filter(move |row| {
std::ops::RangeBounds::contains(&predicate_range, &row.#row_field_ident)
})
}
});

Ok(quote! {
pub fn #fn_name<R>(&self, range: R) -> SelectQueryBuilder<#row_ident,
Expand All @@ -167,9 +228,15 @@ impl InMemoryGenerator {
where
R: #range_bounds
{
#predicate_setup
let read_guard = self.0.data.read_guard();
let rows = self.0.indexes.#field_ident
.range(#range_arg)
.filter_map(|(_, link)| self.0.data.select_non_ghosted(link.0).ok());
.filter_map(move |(_, link)| {
let _read_guard = &read_guard;
self.0.data.select_non_ghosted(link.0).ok()
})
#predicate_filter;

SelectQueryBuilder::new_sorted(rows, #row_fields_ident::#column_pascal)
}
Expand Down
6 changes: 5 additions & 1 deletion codegen/src/generators/persist/queries/select.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,9 +29,13 @@ impl PersistGenerator {
#column_range_type,
#row_fields_ident>
{
let read_guard = self.0.data.read_guard();
let iter = self.0.primary_index.pk_map
.iter()
.filter_map(|(_, link)| self.0.data.select_non_ghosted(link.0).ok());
.filter_map(move |(_, link)| {
let _read_guard = &read_guard;
self.0.data.select_non_ghosted(link.0).ok()
});

SelectQueryBuilder::new(iter)
}
Expand Down
7 changes: 6 additions & 1 deletion codegen/src/generators/persist/table/impls.rs
Original file line number Diff line number Diff line change
Expand Up @@ -196,9 +196,13 @@ impl PersistGenerator {
range.start_bound().map(|v| #primary_key_type::from(v.clone())),
range.end_bound().map(|v| #primary_key_type::from(v.clone())),
);
let read_guard = self.0.data.read_guard();
let rows = self.0.primary_index.pk_map
.range(converted_range)
.filter_map(|(_, link)| self.0.data.select_non_ghosted(link.0).ok());
.filter_map(move |(_, link)| {
let _read_guard = &read_guard;
self.0.data.select_non_ghosted(link.0).ok()
});

#pk_sorted_by
}
Expand Down Expand Up @@ -369,6 +373,7 @@ impl PersistGenerator {

fn gen_table_iter_inner(&self, func: TokenStream) -> TokenStream {
quote! {
let _read_guard = self.0.data.read_guard();
let first = self.0.primary_index.pk_map.iter().next().map(|(k, v)| (k.clone(), v.0));
let Some((mut k, link)) = first else {
return Ok(())
Expand Down
Loading