Skip to content
Merged
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
27 changes: 27 additions & 0 deletions .github/workflows/miri.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
name: Miri

on:
workflow_dispatch:

env:
CARGO_TERM_COLOR: always

jobs:
publication:
runs-on: ubicloud-standard-2
timeout-minutes: 30
steps:
- uses: actions/checkout@v4
- uses: dtolnay/rust-toolchain@nightly
with:
components: miri
- name: Prepare Miri
run: cargo +nightly miri setup
- name: Published-pointer reclamation interleavings
env:
MIRIFLAGS: -Zmiri-many-seeds=0..8
run: cargo +nightly miri test --all-features published_pointer_survives_reclamation_interleavings
- name: Stale-boundary and recovery paths
env:
MIRIFLAGS: -Zmiri-many-seeds=0..8
run: cargo +nightly miri test --all-features attach_
34 changes: 34 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,40 @@ All notable changes to this project will be documented in this file.
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/),
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

## [0.0.12]

### Changed

- Concurrent point lookups route through immutable, chunked,
grace-period-protected topology snapshots instead of acquiring the
structural read lock. Per-node read/write locks let independent readers
proceed without serializing on a node mutex.
- Concurrent CDC `iter_nodes` is replaced by `snapshot_nodes`, which returns
detached `Vec<Node>` values so callers cannot mistake a checkpoint snapshot
for a mutable live node.
- Sequential map exact lookups stop when binary search finds equality instead
of completing a lower-bound search and comparing the result a second time.

### Added

- A split-heavy concurrent regression test verifies that definitive point
reads remain correct while topology generations are repeatedly published.
- `attach_nodes` and `attach_multi_nodes` reconstruct a persisted topology with
one publication instead of publishing once per node.
- A sorted 100,000-key construction benchmark guards the monotonic insert
pattern used by autoincrementing tables.

### Fixed

- Incremental node attachment repairs a deliberately stale final-node route
before adding a later node, preventing false point-lookup misses in the new
node's range. Publication bookkeeping mismatches recover once from the
canonical topology instead of panicking or entering an unbounded collision
chain under the structural write lock.
- Published route chunks keep merge/split headroom, and failed route removals
no longer copy a shared chunk. The bounded point-read fallback releases the
structural lock before invoking caller code.

## [0.0.11]

### Changed
Expand Down
6 changes: 5 additions & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "WorkTablesIndex"
version = "0.0.11"
version = "0.0.12"
edition = "2021"
documentation = "https://docs.rs/WorkTablesIndex/"
repository = "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/pathscale/WorkTablesIndex"
Expand All @@ -13,6 +13,7 @@ readme = "README.md"
[dev-dependencies]
criterion = { version = "0.5", features = ["html_reports"] }
crossbeam-skiplist = "^0.1"
loom = "0.7"
rand = "0.9"
scc = { version = "2.2.5" }
uuid = { version = "1.17.0", features = ["v7"]}
Expand All @@ -33,6 +34,7 @@ parking_lot = { version = "0.12", features = [
"send_guard",
"arc_lock",
], optional = true }
ps-reclaim = { version = "^0.1, >=0.1.3", optional = true }
fastrand = { version = "2", optional = true }
superslice = { version = "1", optional = true }
wt-slice = { version = "0.1", optional = true }
Expand All @@ -42,6 +44,7 @@ default = ["wt-slice-binary-search"]
serde = ["dep:serde"]
concurrent = [
"dep:parking_lot",
"dep:ps-reclaim",
]
cdc = ["concurrent"]
multimap = ["concurrent", "dep:fastrand"]
Expand All @@ -56,6 +59,7 @@ features = ["cdc", "concurrent", "multimap"]
[[bench]]
name = "stdlib"
harness = false
required-features = ["concurrent"]

[[bench]]
name = "concurrent"
Expand Down
37 changes: 36 additions & 1 deletion benches/stdlib.rs
Original file line number Diff line number Diff line change
@@ -1,8 +1,9 @@
use criterion::{criterion_group, criterion_main, Criterion};
use criterion::{criterion_group, criterion_main, BatchSize, BenchmarkId, Criterion, Throughput};
use rand::seq::SliceRandom;
use rand::thread_rng;
use scc::TreeIndex;
use std::hint::black_box;
use std::time::Duration;

fn criterion_benchmark(c: &mut Criterion) {
let n = 100000;
Expand Down Expand Up @@ -42,6 +43,17 @@ fn criterion_benchmark(c: &mut Criterion) {
assert_eq!(indexset.len(), n);
})
});
c.bench_function("concurrent indexset insert sorted 100k", |b| {
b.iter(|| {
let indexset: indexset::concurrent::set::BTreeSet<usize> = indexset::concurrent::set::BTreeSet::new();

for item in 0..n {
black_box(indexset.insert(item));
}

assert_eq!(indexset.len(), n);
})
});
c.bench_function("treeindex insert 100k", |b| {
b.iter(|| {
let treeindex = TreeIndex::new();
Expand All @@ -54,6 +66,29 @@ fn criterion_benchmark(c: &mut Criterion) {
})
});

let mut restore = c.benchmark_group("concurrent indexset restore nodes");
restore.sample_size(20);
restore.warm_up_time(Duration::from_millis(500));
restore.measurement_time(Duration::from_secs(2));
for node_count in [1_000usize, 2_000, 4_000, 8_000] {
let nodes = (0..node_count)
.map(|node| (node * 8..node * 8 + 8).collect::<Vec<_>>())
.collect::<Vec<_>>();
restore.throughput(Throughput::Elements(node_count as u64));
restore.bench_with_input(BenchmarkId::new("bulk", node_count), &nodes, |b, nodes| {
b.iter_batched(
|| nodes.clone(),
|nodes| {
let set = indexset::concurrent::set::BTreeSet::<usize>::with_maximum_node_size(8);
set.attach_nodes(nodes);
black_box(set.node_count());
},
BatchSize::SmallInput,
)
});
}
restore.finish();

let stdlib = std::collections::BTreeSet::from_iter(input.iter());
let indexset = indexset::BTreeSet::from_iter(input.iter());
let concurrent_indexset: indexset::concurrent::set::BTreeSet<usize> = indexset::concurrent::set::BTreeSet::new();
Expand Down
42 changes: 28 additions & 14 deletions src/concurrent/map.rs
Original file line number Diff line number Diff line change
@@ -1,9 +1,6 @@
use std::fmt::{Debug, Display, Formatter};
use std::sync::Arc;
use std::{borrow::Borrow, iter::FusedIterator, ops::RangeBounds};

use parking_lot::Mutex;

use super::set::BTreeSet;
use crate::core::node::NodeLike;
use crate::{cdc::change::ChangeEvent, core::pair::Pair};
Expand Down Expand Up @@ -222,10 +219,27 @@ where
pub fn attach_node(&self, node: Node) {
self.set.attach_node(node)
}
/// Returns iterator over this set's [`Node`]'s.
/// Attaches persisted [`Node`]s with one topology publication.
#[cfg(feature = "cdc")]
pub fn attach_nodes(&self, nodes: impl IntoIterator<Item = Node>) {
self.set.attach_nodes(nodes)
}

/// Returns detached, read-only snapshots of this map's [`Node`]s.
///
/// Callers requiring one coherent logical generation must prevent
/// concurrent mutation while collecting.
#[cfg(feature = "cdc")]
pub fn iter_nodes(&self) -> impl Iterator<Item = Arc<Mutex<Node>>> + '_ {
self.set.index.read().values().cloned().collect::<Vec<_>>().into_iter()
pub fn snapshot_nodes(&self) -> Vec<Node>
where
Node: Clone,
{
self.set
.index
.read()
.values()
.map(|node| (*node.read()).clone())
.collect()
}

/// Copies the exact node boundaries into a pointer-free checkpoint image.
Expand Down Expand Up @@ -274,14 +288,14 @@ where
}

let map = Self::with_maximum_node_size(topology.node_capacity);
for values in topology.nodes {
map.attach_nodes(topology.nodes.into_iter().map(|values| {
let mut node = Node::with_capacity(topology.node_capacity);
for value in values {
let (inserted, _) = NodeLike::insert(&mut node, value);
debug_assert!(inserted, "validated topology contains unique values");
}
map.attach_node(node);
}
node
}));
Ok(map)
}
/// Returns `true` if the map contains a value for the specified key.
Expand Down Expand Up @@ -780,7 +794,7 @@ mod tests {
.index
.read()
.iter()
.map(|(key, node)| (key.clone().key, node.lock_arc().clone()))
.map(|(key, node)| (key.clone().key, node.read_arc().clone()))
.collect::<_>();
assert_eq!(mock_state.nodes, expected_state);
}
Expand Down Expand Up @@ -810,7 +824,7 @@ mod tests {
.index
.read()
.iter()
.map(|(key, node)| (key.clone().key, node.lock_arc().clone()))
.map(|(key, node)| (key.clone().key, node.read_arc().clone()))
.collect::<_>();
assert_eq!(mock_state.nodes, expected_state);
}
Expand Down Expand Up @@ -841,7 +855,7 @@ mod tests {
.index
.read()
.iter()
.map(|(key, node)| (key.clone().key, node.lock_arc().clone()))
.map(|(key, node)| (key.clone().key, node.read_arc().clone()))
.collect::<_>();
assert_eq!(mock_state.nodes, expected_state);
}
Expand Down Expand Up @@ -874,7 +888,7 @@ mod tests {
.index
.read()
.iter()
.map(|(key, node)| (key.clone().key, node.lock_arc().clone()))
.map(|(key, node)| (key.clone().key, node.read_arc().clone()))
.collect::<_>();
assert_eq!(mock_state.nodes, expected_state);
}
Expand Down Expand Up @@ -941,7 +955,7 @@ mod tests {
.index
.read()
.iter()
.map(|(key, node)| (key.clone().key, node.lock_arc().clone()))
.map(|(key, node)| (key.clone().key, node.read_arc().clone()))
.collect::<_>();
assert_eq!(mock_state.nodes, expected_state);
}
Expand Down
26 changes: 20 additions & 6 deletions src/concurrent/multimap.rs
Original file line number Diff line number Diff line change
@@ -1,14 +1,11 @@
use std::fmt::Debug;
use std::marker::PhantomData;
use std::sync::Arc;
use std::{
borrow::Borrow,
iter::FusedIterator,
ops::{Bound, RangeBounds},
};

use parking_lot::Mutex;

use crate::core::node::NodeLike;
use crate::{
cdc::change::ChangeEvent,
Expand Down Expand Up @@ -210,10 +207,27 @@ where
pub fn attach_multi_node(&self, node: Node) {
self.set.attach_node(node)
}
/// Returns iterator over this multiset's [`Node`]'s.
/// Attaches persisted [`Node`]s with one topology publication.
#[cfg(feature = "cdc")]
pub fn iter_nodes(&self) -> impl Iterator<Item = Arc<Mutex<Node>>> + '_ {
self.set.index.read().values().cloned().collect::<Vec<_>>().into_iter()
pub fn attach_multi_nodes(&self, nodes: impl IntoIterator<Item = Node>) {
self.set.attach_nodes(nodes)
}

/// Returns detached, read-only snapshots of this multimap's [`Node`]s.
///
/// Callers requiring one coherent logical generation must prevent
/// concurrent mutation while collecting.
#[cfg(feature = "cdc")]
pub fn snapshot_nodes(&self) -> Vec<Node>
where
Node: Clone,
{
self.set
.index
.read()
.values()
.map(|node| (*node.read()).clone())
.collect()
}
/// Returns `true` if the map contains at least one occurance of the specified key.
///
Expand Down
Loading
Loading