From 7cc215793b47a80d90da178929eb99b73942aa95 Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 5 Sep 2026 17:53:03 +0700 Subject: [PATCH 1/4] Remove shared reader bottlenecks from WTI lookups --- CHANGELOG.md | 18 ++ Cargo.toml | 4 +- src/concurrent/map.rs | 29 ++- src/concurrent/multimap.rs | 19 +- src/concurrent/operation.rs | 16 +- src/concurrent/ref.rs | 4 +- src/concurrent/set.rs | 392 ++++++++++++++++++++++++++++-------- src/core/node.rs | 43 ++++ src/lib.rs | 18 +- 9 files changed, 422 insertions(+), 121 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index a1d0be4..29e0f0c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,24 @@ 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, 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` retains its mutex-shaped API but returns detached + node snapshots. Mutating a returned node no longer mutates the live index. +- 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. + ## [0.0.11] ### Changed diff --git a/Cargo.toml b/Cargo.toml index bc2080c..a06146e 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -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" @@ -33,6 +33,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 } @@ -42,6 +43,7 @@ default = ["wt-slice-binary-search"] serde = ["dep:serde"] concurrent = [ "dep:parking_lot", + "dep:ps-reclaim", ] cdc = ["concurrent"] multimap = ["concurrent", "dep:fastrand"] diff --git a/src/concurrent/map.rs b/src/concurrent/map.rs index 5f0872a..f6d400c 100644 --- a/src/concurrent/map.rs +++ b/src/concurrent/map.rs @@ -222,10 +222,23 @@ where pub fn attach_node(&self, node: Node) { self.set.attach_node(node) } - /// Returns iterator over this set's [`Node`]'s. + /// Returns detached snapshots of this map's [`Node`]s. + /// + /// The returned mutexes preserve the checkpoint API's existing shape, but + /// mutating them does not mutate this map. Callers requiring one coherent + /// logical generation must prevent concurrent mutation while collecting. #[cfg(feature = "cdc")] - pub fn iter_nodes(&self) -> impl Iterator>> + '_ { - self.set.index.read().values().cloned().collect::>().into_iter() + pub fn iter_nodes(&self) -> impl Iterator>> + '_ + where + Node: Clone, + { + self.set + .index + .read() + .values() + .map(|node| Arc::new(Mutex::new((*node.read()).clone()))) + .collect::>() + .into_iter() } /// Copies the exact node boundaries into a pointer-free checkpoint image. @@ -780,7 +793,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); } @@ -810,7 +823,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); } @@ -841,7 +854,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); } @@ -874,7 +887,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); } @@ -941,7 +954,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); } diff --git a/src/concurrent/multimap.rs b/src/concurrent/multimap.rs index 360c1f0..c115aeb 100644 --- a/src/concurrent/multimap.rs +++ b/src/concurrent/multimap.rs @@ -210,10 +210,23 @@ where pub fn attach_multi_node(&self, node: Node) { self.set.attach_node(node) } - /// Returns iterator over this multiset's [`Node`]'s. + /// Returns detached snapshots of this multimap's [`Node`]s. + /// + /// The returned mutexes preserve the checkpoint API's existing shape, but + /// mutating them does not mutate this map. Callers requiring one coherent + /// logical generation must prevent concurrent mutation while collecting. #[cfg(feature = "cdc")] - pub fn iter_nodes(&self) -> impl Iterator>> + '_ { - self.set.index.read().values().cloned().collect::>().into_iter() + pub fn iter_nodes(&self) -> impl Iterator>> + '_ + where + Node: Clone, + { + self.set + .index + .read() + .values() + .map(|node| Arc::new(Mutex::new((*node.read()).clone()))) + .collect::>() + .into_iter() } /// Returns `true` if the map contains at least one occurance of the specified key. /// diff --git a/src/concurrent/operation.rs b/src/concurrent/operation.rs index 067b4ee..9c87f13 100644 --- a/src/concurrent/operation.rs +++ b/src/concurrent/operation.rs @@ -1,14 +1,14 @@ use std::fmt::Debug; use std::sync::Arc; -use parking_lot::Mutex; +use parking_lot::RwLock; use std::collections::BTreeMap; use crate::cdc::change::ChangeEventUnassigned; use crate::core::node::NodeLike; -type OldVersion = Arc>; -type CurrentVersion = Arc>; +type OldVersion = Arc>; +type CurrentVersion = Arc>; pub enum Operation> { Split(OldVersion, T, T), @@ -27,12 +27,12 @@ where // value in place: see `MultiPairLike::adopt_stored_identity`. pub fn commit( self, - index: &mut BTreeMap>>, + index: &mut BTreeMap>>, adopt: fn(&T, &mut T), ) -> Result<(Option, Vec>), ()> { match self { Operation::Split(old_node, old_max, value) => { - let mut guard = old_node.lock_arc(); + let mut guard = old_node.write_arc(); if let Some(entry) = index.get(&old_max) { if Arc::ptr_eq(entry, &old_node) { // The node was drained by a concurrent remove after @@ -131,7 +131,7 @@ where cdc.push(value_insertion); } } - let new_node = Arc::new(Mutex::new(new_vec)); + let new_node = Arc::new(RwLock::new(new_vec)); index.insert(max, new_node); } @@ -143,7 +143,7 @@ where Err(()) } Operation::UpdateMax(node, old_max) => { - let guard = node.lock_arc(); + let guard = node.write_arc(); if let Some(entry) = index.get(&old_max) { if Arc::ptr_eq(entry, &node) { let mut cdc = vec![]; @@ -188,7 +188,7 @@ where Err(()) } Operation::MakeUnreachable(node, old_max) => { - let guard = node.lock_arc(); + let guard = node.write_arc(); if let Some(entry) = index.get(&old_max) { if Arc::ptr_eq(entry, &node) { return match guard.max() { diff --git a/src/concurrent/ref.rs b/src/concurrent/ref.rs index d27178d..0bf31d4 100644 --- a/src/concurrent/ref.rs +++ b/src/concurrent/ref.rs @@ -1,9 +1,9 @@ use crate::core::node::NodeLike; -use parking_lot::{ArcMutexGuard, RawMutex}; +use parking_lot::{ArcRwLockReadGuard, RawRwLock}; use std::marker::PhantomData; pub struct Ref + Send> { - pub(super) node_guard: ArcMutexGuard, + pub(super) node_guard: ArcRwLockReadGuard, pub(super) position: usize, pub(super) phantom_data: PhantomData, } diff --git a/src/concurrent/set.rs b/src/concurrent/set.rs index 53c83f2..f29cfe8 100644 --- a/src/concurrent/set.rs +++ b/src/concurrent/set.rs @@ -1,11 +1,10 @@ -use parking_lot::{ArcMutexGuard, Mutex, RawMutex, RwLock}; +use parking_lot::{ArcRwLockReadGuard, ArcRwLockWriteGuard, RawRwLock, RwLock, RwLockReadGuard, RwLockWriteGuard}; use std::collections::BTreeMap; use std::fmt::Debug; use std::iter::FusedIterator; use std::marker::PhantomData; use std::ops::{Bound, RangeBounds}; -#[cfg(feature = "cdc")] -use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::atomic::{AtomicPtr, AtomicU64, Ordering}; use std::{borrow::Borrow, sync::Arc}; use crate::cdc::change::ChangeEvent; @@ -15,13 +14,155 @@ use crate::core::node::*; use super::r#ref::Ref; -// Most point reads acquire the node mutex immediately while the structural -// mapping is pinned. If a node is genuinely contended, wait without holding a -// structural shard and retry. A bounded fallback preserves progress if the -// node keeps being reacquired between the wait and the next stable attempt. -const STABLE_READ_BLOCKING_FALLBACK_AFTER: usize = 2; const ROOT_PUBLICATION_SPIN_LIMIT: usize = 16; +type NodeIndex = BTreeMap>>; + +struct RetiredIndex(*mut NodeIndex); + +// SAFETY: the pointer is uniquely owned after it has been swapped out of the +// publication slot, and this wrapper exposes no access to the map. Its only +// operation is destruction after the grace period. Dropping the keys and, for +// the final Arc, moving each node into its destructor on that thread is valid +// when both stored types are `Send`; sharing `Node` there is not required. +unsafe impl Send for RetiredIndex {} + +impl Drop for RetiredIndex { + fn drop(&mut self) { + // SAFETY: this wrapper is created exactly once for a pointer returned + // by `Box::into_raw`, after that pointer has been atomically unlinked. + unsafe { drop(Box::from_raw(self.0)) } + } +} + +#[derive(Debug)] +struct PublishedIndex { + current: AtomicPtr>, + domain: ps_reclaim::Domain, +} + +impl PublishedIndex { + fn new() -> Self { + Self { + current: AtomicPtr::new(Box::into_raw(Box::new(BTreeMap::new()))), + domain: ps_reclaim::Domain::new(), + } + } +} + +impl PublishedIndex +where + T: Ord + Clone + Send + 'static, + Node: Send + 'static, +{ + fn publish(&self, index: &NodeIndex) { + let replacement = Box::into_raw(Box::new(index.clone())); + let retired = self.current.swap(replacement, Ordering::AcqRel); + // Keep the pointer's provenance intact while transferring its unique + // ownership to the retirement callback. + let retired = RetiredIndex(retired); + self.domain.retire(move || drop(retired)); + // Topology publication is infrequent (normally one node boundary per + // 1,024 inserts), so reclaim one expired snapshot on the writer path. + self.domain.advance_up_to(1); + } +} + +impl Drop for PublishedIndex { + fn drop(&mut self) { + let current = *self.current.get_mut(); + // SAFETY: exclusive access proves no reader can load `current`, and it + // is the one still-linked allocation created by `Box::into_raw`. + unsafe { drop(Box::from_raw(current)) } + } +} + +#[derive(Debug)] +pub(crate) struct Topology { + index: RwLock>, + published: PublishedIndex, + // Even values are stable publications; odd values mean a writer may have + // changed node contents or routing but has not published the new route. + generation: AtomicU64, +} + +impl Topology { + fn new() -> Self { + Self { + index: RwLock::new(BTreeMap::new()), + published: PublishedIndex::new(), + generation: AtomicU64::new(0), + } + } + + #[inline] + pub(crate) fn read(&self) -> RwLockReadGuard<'_, NodeIndex> { + self.index.read() + } +} + +impl Topology +where + T: Ord + Clone + Send + 'static, + Node: Send + 'static, +{ + #[inline] + fn write(&self) -> TopologyWriteGuard<'_, T, Node> { + let index = self.index.write(); + self.generation.fetch_add(1, Ordering::AcqRel); + TopologyWriteGuard { topology: self, index } + } + + #[inline] + fn try_write(&self) -> Option> { + let index = self.index.try_write()?; + self.generation.fetch_add(1, Ordering::AcqRel); + Some(TopologyWriteGuard { topology: self, index }) + } +} + +struct TopologyWriteGuard<'a, T, Node> +where + T: Ord + Clone + Send + 'static, + Node: Send + 'static, +{ + topology: &'a Topology, + index: RwLockWriteGuard<'a, NodeIndex>, +} + +impl std::ops::Deref for TopologyWriteGuard<'_, T, Node> +where + T: Ord + Clone + Send + 'static, + Node: Send + 'static, +{ + type Target = NodeIndex; + + fn deref(&self) -> &Self::Target { + &self.index + } +} + +impl std::ops::DerefMut for TopologyWriteGuard<'_, T, Node> +where + T: Ord + Clone + Send + 'static, + Node: Send + 'static, +{ + fn deref_mut(&mut self) -> &mut Self::Target { + &mut self.index + } +} + +impl Drop for TopologyWriteGuard<'_, T, Node> +where + T: Ord + Clone + Send + 'static, + Node: Send + 'static, +{ + fn drop(&mut self) { + self.topology.published.publish(&self.index); + self.topology.generation.fetch_add(1, Ordering::Release); + } +} + // Default identity-adoption hook for replace-on-equality: plain sets and maps // have no hidden ordering state to carry over. See // `MultiPairLike::adopt_stored_identity`. @@ -144,14 +285,12 @@ where T: Ord + Clone + 'static, Node: NodeLike, { - // The old representation put a concurrent Crossbeam skip-list behind a - // second structural read/write lock. Every topology read already held the - // outer lock and every topology mutation held it exclusively, so the - // skip-list's epoch reclamation and atomics could not add concurrency. - // Keeping the ordered map inside the one structural lock removes that - // redundant reclamation domain. Node contents retain their independent - // mutexes and remain concurrently mutable. - pub(crate) index: RwLock>>>, + // Writers maintain the canonical ordered topology under one structural + // lock and publish immutable snapshots for point reads. The read path is + // therefore free of a shared reader-count cache line. Node contents use + // independent read/write locks, so readers routed to one node may proceed + // concurrently while mutations retain exclusive node access. + pub(crate) index: Topology, node_capacity: usize, // Ordinary set/map keys satisfy Borrow's ordering contract and retain a // logarithmic BTreeMap route. Multimap entries intentionally borrow only @@ -166,7 +305,7 @@ where impl> Default for BTreeSet { fn default() -> Self { Self { - index: RwLock::new(BTreeMap::new()), + index: Topology::new(), node_capacity: DEFAULT_INNER_SIZE, borrow_order_matches: true, #[cfg(feature = "cdc")] @@ -194,7 +333,7 @@ where /// let set: BTreeSet = BTreeSet::with_maximum_node_size(128); pub fn with_maximum_node_size(node_capacity: usize) -> Self { Self { - index: RwLock::new(BTreeMap::new()), + index: Topology::new(), node_capacity, borrow_order_matches: true, #[cfg(feature = "cdc")] @@ -210,7 +349,7 @@ where .max() .cloned() .expect("node should contain at least one value to be correct node"); - self.index.write().insert(node_id, Arc::new(Mutex::new(node))); + self.index.write().insert(node_id, Arc::new(RwLock::new(node))); } #[cfg(feature = "cdc")] @@ -218,7 +357,7 @@ where let index = self.index.read(); let nodes = index .values() - .map(|node| node.lock().iter().cloned().collect()) + .map(|node| node.read().iter().cloned().collect()) .collect(); (self.node_capacity, nodes) } @@ -230,7 +369,7 @@ where &self, value: T, adopt: fn(&T, &mut T), - ) -> Result<(Option, Vec>), (ArcMutexGuard, usize, T)> { + ) -> Result<(Option, Vec>), (ArcRwLockWriteGuard, usize, T)> { loop { let mut cdc = vec![]; let index = self.index.read(); @@ -275,14 +414,14 @@ where cdc.push(node_insertion); } - index.insert(value, Arc::new(Mutex::new(first_node))); + index.insert(value, Arc::new(RwLock::new(first_node))); return Ok((None, cdc)); } } }; - let mut node_guard = target_node_entry.1.clone().lock_arc(); + let mut node_guard = target_node_entry.1.clone().write_arc(); #[allow(unused_assignments)] let mut operation = None; @@ -462,7 +601,7 @@ where pub(crate) fn put_checked( &self, value: T, - ) -> Result<(Option, Vec>), (ArcMutexGuard, usize, T)> { + ) -> Result<(Option, Vec>), (ArcRwLockWriteGuard, usize, T)> { self.put_checked_inner::(value, no_identity_adoption) } @@ -479,7 +618,7 @@ where pub(crate) fn put_cdc_checked( &self, value: T, - ) -> Result<(Option, Vec>), (ArcMutexGuard, usize, T)> { + ) -> Result<(Option, Vec>), (ArcRwLockWriteGuard, usize, T)> { self.put_checked_inner::(value, no_identity_adoption) } @@ -525,7 +664,7 @@ where first_for_borrowed_bound(&index, Bound::Included(value), self.borrow_order_matches) .or_else(|| index.last_key_value()) { - let mut node_guard = target_node_entry.1.clone().lock_arc(); + let mut node_guard = target_node_entry.1.clone().write_arc(); let old_max = node_guard.max().cloned(); let deleted = NodeLike::delete(&mut *node_guard, value); if deleted.is_none() { @@ -627,7 +766,7 @@ where // the multimap paths use this, and it relies on NodeLike::delete_at (also // multimap-gated), so gate the whole family to avoid an unconditional break. #[inline(always)] - fn lock_node_for_value_optimistic(&self, value: &Q) -> Option> + fn lock_node_for_value_optimistic(&self, value: &Q) -> Option> where T: Borrow, Q: Ord + ?Sized, @@ -642,58 +781,58 @@ where .or_else(|| index.first_key_value().map(|(_, node)| node.clone())), } }?; - Some(node.lock_arc()) + Some(node.read_arc()) } /// Locates and locks the node whose structural range owns `value`. /// - /// The common uncontended path acquires the node with `try_lock_arc` while - /// holding the topology read lock, making both hits and misses definitive with - /// one structural lookup. On contention, the structural guard is released - /// before waiting so a long-lived node reference cannot convoy unrelated - /// structural writers. After repeated contention, the documented - /// topology-to-node lock order is used as a bounded progress fallback. + /// Readers route through an immutable published topology and validate its + /// generation after locking the node. A concurrent structural change makes + /// the read retry, so hits and misses remain definitive without updating a + /// shared reader-count cache line. #[inline(always)] - fn lock_node_for_value(&self, value: &Q) -> Option> + fn lock_node_for_value(&self, value: &Q) -> Option> where T: Borrow, Q: Ord + ?Sized, { - let mut contentions = 0; - loop { - let index = self.index.read(); - let node = match first_for_borrowed_bound(&index, Bound::Included(value), self.borrow_order_matches) { - Some((_, node)) => node.clone(), + let generation = self.index.generation.load(Ordering::Acquire); + if !generation.is_multiple_of(2) { + std::hint::spin_loop(); + continue; + } + + let pin = self.index.published.domain.pin(); + let snapshot = self.index.published.current.load(Ordering::Acquire); + // SAFETY: `snapshot` was loaded after `pin`, and the publication + // domain cannot reclaim it until `pin` is dropped. + let index = unsafe { &*snapshot }; + let node = match first_for_borrowed_bound(index, Bound::Included(value), self.borrow_order_matches) { + Some((_, node)) => Some(node.clone()), None => index .last_key_value() .map(|(_, node)| node.clone()) - .or_else(|| index.first_key_value().map(|(_, node)| node.clone()))?, + .or_else(|| index.first_key_value().map(|(_, node)| node.clone())), + }; + let Some(node) = node else { + if self.index.generation.load(Ordering::Acquire) == generation { + return None; + } + continue; }; - if let Some(node_guard) = node.try_lock_arc() { - drop(index); - return Some(node_guard); - } - - contentions += 1; - if contentions >= STABLE_READ_BLOCKING_FALLBACK_AFTER { - let node_guard = node.lock_arc(); - drop(index); + let node_guard = node.read_arc(); + if self.index.generation.load(Ordering::Acquire) == generation { + drop(pin); return Some(node_guard); } - - drop(index); - // Wait for the observed holder without pinning the structural - // mapping, then retry so the returned guard always corresponds to - // a mapping observed under the topology read lock. - drop(node.lock_arc()); } } #[inline(always)] fn get_with_guard( - node_guard: ArcMutexGuard, + node_guard: ArcRwLockReadGuard, value: &Q, read: impl FnOnce(&T) -> R, ) -> Option @@ -711,7 +850,39 @@ where T: Borrow, Q: Ord + ?Sized, { - Self::get_with_guard(self.lock_node_for_value(value)?, value, read) + loop { + let generation = self.index.generation.load(Ordering::Acquire); + if !generation.is_multiple_of(2) { + std::hint::spin_loop(); + continue; + } + + let _pin = self.index.published.domain.pin(); + let snapshot = self.index.published.current.load(Ordering::Acquire); + // SAFETY: `snapshot` was loaded after `pin`, and the publication + // domain cannot reclaim it until `pin` is dropped. + let index = unsafe { &*snapshot }; + let node = first_for_borrowed_bound(index, Bound::Included(value), self.borrow_order_matches) + .or_else(|| index.last_key_value()) + .or_else(|| index.first_key_value()) + .map(|(_, node)| node); + let Some(node) = node else { + if self.index.generation.load(Ordering::Acquire) == generation { + return None; + } + continue; + }; + + // Borrow the Arc from the protected snapshot: unlike + // `lock_node_for_value`, this owned-result path does not need an + // Arc clone or its shared refcount update. + let node_guard = node.read(); + if self.index.generation.load(Ordering::Acquire) != generation { + continue; + } + let position = node_guard.try_select(value)?; + return node_guard.get_ith(position).map(read); + } } #[inline(always)] @@ -743,8 +914,7 @@ where T: Borrow, Q: Ord + ?Sized, { - self.lock_node_for_value(value) - .is_some_and(|node_guard| node_guard.contains(value)) + self.get_with(value, |_| ()).is_some() } pub fn get<'a, Q>(&'a self, value: &'a Q) -> Option> where @@ -767,17 +937,17 @@ where } pub fn len(&self) -> usize { - self.index.read().values().map(|node| node.lock().len()).sum() + self.index.read().values().map(|node| node.read().len()).sum() } pub fn is_empty(&self) -> bool { - self.index.read().values().all(|node| node.lock().is_empty()) + self.index.read().values().all(|node| node.read().is_empty()) } pub fn capacity(&self) -> usize { self.index .read() .values() .map(|node| { - let guard = node.lock(); + let guard = node.read(); guard.capacity() }) .sum() @@ -840,16 +1010,16 @@ where // Identity of the node the last batch in each direction was cloned from, // so the next install can step past it when the cursor lookup lands on // it again (its entry key can sit past every element it still holds). - exhausted_front_node: Option>>, - exhausted_back_node: Option>>, + exhausted_front_node: Option>>, + exhausted_back_node: Option>>, // The node a direction is partway through, and how many of its elements it // has taken. A batch that stops short of a node's end must resume inside // that node, and must never take less than it already has: the rank-based // skip alone cannot guarantee that, because a repositioned node can leave // the cursor ranking below elements already yielded. Recording the count // makes forward progress structural rather than incidental. - front_partial: Option<(Arc>, usize)>, - back_partial: Option<(Arc>, usize)>, + front_partial: Option<(Arc>, usize)>, + back_partial: Option<(Arc>, usize)>, // How many elements the next batch may clone, doubling per install. front_batch_limit: usize, back_batch_limit: usize, @@ -941,7 +1111,7 @@ where return false; }; let node = entry.clone(); - let guard = node.lock_arc(); + let guard = node.read_arc(); drop(index); let rank_skip = self @@ -1023,7 +1193,7 @@ where return false; }; let node = entry.clone(); - let guard = node.lock_arc(); + let guard = node.read_arc(); drop(index); let truncate = self @@ -1203,7 +1373,7 @@ where let current_front_entry = first_for_borrowed_bound(&index, start_bound, btree.borrow_order_matches); let front_value = if let Some((front_key, front_node)) = current_front_entry { - let front_guard = front_node.clone().lock_arc(); + let front_guard = front_node.clone().read_arc(); let rank = match start_bound { Bound::Included(v) => front_guard.rank(Bound::Included(v), true), Bound::Excluded(v) => front_guard.rank(Bound::Excluded(v), true), @@ -1222,7 +1392,7 @@ where // Never hold two node locks here. drop(front_guard); if let Some((_, pre_front_node)) = index.range::(..front_key).next_back() { - let pre_front_guard = pre_front_node.clone().lock_arc(); + let pre_front_guard = pre_front_node.clone().read_arc(); pre_front_guard.iter().last().cloned() } else { None @@ -1240,7 +1410,7 @@ where }; let back_value = if let Some((back_key, back_node)) = current_back_entry { - let back_guard = back_node.clone().lock_arc(); + let back_guard = back_node.clone().read_arc(); let rank = match end_bound { Bound::Included(v) => back_guard.rank(Bound::Included(v), false), Bound::Excluded(v) => back_guard.rank(Bound::Excluded(v), false), @@ -1258,7 +1428,7 @@ where .range::((Bound::Excluded(back_key), Bound::Unbounded)) .next() { - let next_back_guard = next_back_node.clone().lock_arc(); + let next_back_guard = next_back_node.clone().read_arc(); next_back_guard.iter().next().cloned() } else { None @@ -1273,7 +1443,7 @@ where if start_bound != Bound::Unbounded || end_bound != Bound::Unbounded { if let Some(max) = index .last_key_value() - .and_then(|(_, node)| node.clone().lock_arc().max().cloned()) + .and_then(|(_, node)| node.clone().read_arc().max().cloned()) { if let Bound::Included(v) = start_bound { if v > max.borrow() { @@ -1288,7 +1458,7 @@ where if let Some(min) = index .first_key_value() - .and_then(|(_, node)| node.clone().lock_arc().min().cloned()) + .and_then(|(_, node)| node.clone().read_arc().min().cloned()) { if let Bound::Included(v) = end_bound { if v < min.borrow() { @@ -1492,7 +1662,7 @@ where if Arc::ptr_eq(&front_node, &back_node) { // The whole range lives in one node. - let mut guard = front_node.clone().lock_arc(); + let mut guard = front_node.clone().write_arc(); let front_position = guard.rank(start_bound, true).map_or(0, |last| last + 1); let back_position = removed_prefix_len(&guard); if back_position <= front_position { @@ -1512,8 +1682,8 @@ where return; } - let mut front_guard = front_node.clone().lock_arc(); - let mut back_guard = back_node.clone().lock_arc(); + let mut front_guard = front_node.clone().write_arc(); + let mut back_guard = back_node.clone().write_arc(); let front_position = front_guard.rank(start_bound, true).map_or(0, |last| last + 1); let back_position = removed_prefix_len(&back_guard); @@ -1526,7 +1696,7 @@ where let node = index .remove::(&key) .expect("middle key was collected under the write lock"); - let mut removed_node = node.lock_arc(); + let mut removed_node = node.write_arc(); detached_nodes.push(std::mem::take(&mut *removed_node)); } @@ -1617,6 +1787,50 @@ mod tests { assert_eq!(set.iter().collect::>(), expected); } + #[test] + fn published_point_reads_remain_definitive_across_splits() { + use std::sync::atomic::{AtomicBool, Ordering}; + + const STABLE_KEYS: usize = 256; + const FINAL_KEYS: usize = 4_096; + const READERS: usize = 4; + + let set = Arc::new(BTreeSet::::with_maximum_node_size(8)); + for key in 0..STABLE_KEYS { + set.insert(key); + } + + let start = Arc::new(Barrier::new(READERS + 1)); + let done = Arc::new(AtomicBool::new(false)); + let readers = (0..READERS) + .map(|reader| { + let set = Arc::clone(&set); + let start = Arc::clone(&start); + let done = Arc::clone(&done); + thread::spawn(move || { + start.wait(); + let mut probe = reader; + while !done.load(Ordering::Acquire) { + let key = probe % STABLE_KEYS; + assert_eq!(set.get_with(&key, |value| *value), Some(key)); + probe += READERS; + } + }) + }) + .collect::>(); + + start.wait(); + for key in STABLE_KEYS..FINAL_KEYS { + set.insert(key); + } + done.store(true, Ordering::Release); + + for reader in readers { + reader.join().unwrap(); + } + assert_eq!(set.len(), FINAL_KEYS); + } + #[test] fn mixed_structural_and_node_lock_paths_complete_without_deadlock() { const THREADS: usize = 8; @@ -1753,7 +1967,7 @@ mod tests { set.index .read() .values() - .flat_map(|node| node.lock().iter().cloned().collect::>()) + .flat_map(|node| node.read().iter().cloned().collect::>()) .collect::>() .symmetric_difference(&inserted_values) .collect::>() @@ -2079,12 +2293,12 @@ mod tests { .unwrap() .1 .clone(); - let detached_values = detached.lock().iter().copied().collect::>(); + let detached_values = detached.read().iter().copied().collect::>(); assert!(detached_values.iter().all(|value| (2..30).contains(value))); set.remove_range(2..30); - assert!(detached.lock().is_empty()); + assert!(detached.read().is_empty()); assert!(detached_values.iter().all(|value| !set.contains(value))); } @@ -2101,7 +2315,7 @@ mod tests { // reach it the same way. { let node = set.index.read().last_key_value().expect("node must exist").1.clone(); - let mut guard = node.lock(); + let mut guard = node.write(); NodeLike::insert(&mut *guard, 5u64); } @@ -2117,7 +2331,7 @@ mod tests { fn drain_node_with_pending_unlink(set: &BTreeSet, values: &[u64], stale_key: u64) -> Operation> { let node = set.index.read().last_key_value().expect("node must exist").1.clone(); { - let mut guard = node.lock(); + let mut guard = node.write(); for value in values { NodeLike::delete(&mut *guard, value).expect("seeded value must be present"); } @@ -2136,7 +2350,7 @@ mod tests { let pending_split = Operation::Split(node.clone(), 30u64, 15u64); // ...then a concurrent remove drains the node before the commit. { - let mut guard = node.lock(); + let mut guard = node.write(); for seeded in [10u64, 20, 30] { NodeLike::delete(&mut *guard, &seeded).expect("seeded value must be present"); } @@ -2169,7 +2383,7 @@ mod tests { let node = set.index.read().last_key_value().expect("node must exist").1.clone(); let pending_split = Operation::Split(node.clone(), 30u64, 15u64); { - let mut guard = node.lock(); + let mut guard = node.write(); for seeded in [10u64, 20, 30] { NodeLike::delete(&mut *guard, &seeded).expect("seeded value must be present"); } @@ -2227,7 +2441,7 @@ mod tests { // routed (empty) node under the node lock; the UpdateMax repair // has not committed yet. { - let mut guard = node.lock(); + let mut guard = node.write(); NodeLike::insert(&mut *guard, value); } let pending_repair = Operation::UpdateMax(node.clone(), 30u64); @@ -2277,7 +2491,7 @@ mod tests { let mut attempts = 0; while set.remove(&key).is_none() { let index = set.index.read(); - let present = index.values().any(|node| node.lock().contains(&key)); + let present = index.values().any(|node| node.read().contains(&key)); drop(index); assert!(present, "acknowledged insert of {key} was lost"); attempts += 1; @@ -2686,7 +2900,7 @@ mod tests { set.attach_node(vec![30u64, 40]); { let node = set.index.read().first_key_value().expect("fixture node").1.clone(); - let mut guard = node.lock(); + let mut guard = node.write(); NodeLike::insert(&mut *guard, 35u64); } diff --git a/src/core/node.rs b/src/core/node.rs index 13cfc2e..4d3dc9b 100644 --- a/src/core/node.rs +++ b/src/core/node.rs @@ -70,6 +70,11 @@ mod search_backend { { haystack.binary_search_by(|candidate| candidate.borrow().cmp(needle)) } + + #[inline] + pub(crate) fn search_by(haystack: &[T], compare: impl FnMut(&T) -> core::cmp::Ordering) -> Result { + haystack.binary_search_by(compare) + } } #[cfg(all( @@ -99,6 +104,18 @@ mod search_backend { _ => Err(index), } } + + #[inline] + pub(crate) fn search_by( + haystack: &[T], + mut compare: impl FnMut(&T) -> core::cmp::Ordering, + ) -> Result { + let index = haystack.lower_bound_by(&mut compare); + match haystack.get(index) { + Some(candidate) if compare(candidate).is_eq() => Ok(index), + _ => Err(index), + } + } } #[cfg(all( @@ -120,6 +137,11 @@ mod search_backend { { haystack.exact_search_by(|candidate| candidate.borrow().cmp(needle)) } + + #[inline] + pub(crate) fn search_by(haystack: &[T], compare: impl FnMut(&T) -> core::cmp::Ordering) -> Result { + haystack.exact_search_by(compare) + } } #[cfg(any( @@ -167,6 +189,26 @@ mod search_backend { Err(i) } + + #[inline] + pub(crate) fn search_by(haystack: &[T], mut compare: impl FnMut(&T) -> Ordering) -> Result { + let mut right = haystack.len(); + let mut left = 0; + + while left != right { + let middle = (left + right) >> 1; + // SAFETY: `left < right <= haystack.len()` makes `middle` a valid + // element index, and both branches preserve the bounds. + let candidate = unsafe { haystack.get_unchecked(middle) }; + match compare(candidate) { + Ordering::Equal => return Ok(middle), + Ordering::Less => left = middle + 1, + Ordering::Greater => right = middle, + } + } + + Err(left) + } } // Search backend precedence is deterministic when features are composed: @@ -174,6 +216,7 @@ mod search_backend { // custom implementation is the compatibility fallback. Each backend's cfg // selects its implementation and test name together, preventing drift. use search_backend::search; +pub(crate) use search_backend::search_by; #[inline] fn compute_positions_to_skip(haystack: &[T], bound: std::ops::Bound<&Q>, forward: bool) -> Option diff --git a/src/lib.rs b/src/lib.rs index 414d077..0ef3a0c 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -2162,16 +2162,14 @@ where K: Borrow + Ord, Q: Ord + ?Sized, { - let (node_idx, position_within_node) = self.set.locate_value_cmp(|item: &Pair| item.key.borrow() < key); - if let Some(candidate_node) = self.set.inner.get(node_idx) { - if let Some(candidate_value) = candidate_node.get(position_within_node) { - if candidate_value.key.borrow() == key { - return Some((&candidate_value.key, &candidate_value.value)); - } - } - } - - None + let node_idx = self.set.locate_node_cmp(|item: &Pair| item.key.borrow() < key); + let candidate_node = self.set.inner.get(node_idx)?; + let position = crate::core::node::search_by(candidate_node, |candidate| { + >::borrow(&candidate.key).cmp(key) + }) + .ok()?; + let candidate = &candidate_node[position]; + Some((&candidate.key, &candidate.value)) } /// Returns a mutable reference to the value corresponding to the key. /// From 7d448df6d9e773c45a9a5bc46b6a899200e76459 Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 5 Sep 2026 19:55:42 +0700 Subject: [PATCH 2/4] Finish scalable topology publication --- CHANGELOG.md | 17 +- Cargo.toml | 2 + benches/stdlib.rs | 37 +- src/concurrent/map.rs | 29 +- src/concurrent/multimap.rs | 23 +- src/concurrent/operation.rs | 18 +- src/concurrent/ref.rs | 6 + src/concurrent/set.rs | 739 +++++++++++++++++++++++++++++++++--- src/core/node.rs | 67 ++-- src/lib.rs | 4 +- tests/loom_publication.rs | 63 +++ 11 files changed, 876 insertions(+), 129 deletions(-) create mode 100644 tests/loom_publication.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index 29e0f0c..2054f46 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,12 +9,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Changed -- Concurrent point lookups route through immutable, 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` retains its mutex-shaped API but returns detached - node snapshots. Mutating a returned node no longer mutates the live index. +- 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` 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. @@ -22,6 +23,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - 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. ## [0.0.11] diff --git a/Cargo.toml b/Cargo.toml index a06146e..4a27a33 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -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"]} @@ -58,6 +59,7 @@ features = ["cdc", "concurrent", "multimap"] [[bench]] name = "stdlib" harness = false +required-features = ["concurrent"] [[bench]] name = "concurrent" diff --git a/benches/stdlib.rs b/benches/stdlib.rs index d5486b7..fc9aed0 100644 --- a/benches/stdlib.rs +++ b/benches/stdlib.rs @@ -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; @@ -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 = 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(); @@ -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::>()) + .collect::>(); + 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::::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 = indexset::concurrent::set::BTreeSet::new(); diff --git a/src/concurrent/map.rs b/src/concurrent/map.rs index f6d400c..90f92c7 100644 --- a/src/concurrent/map.rs +++ b/src/concurrent/map.rs @@ -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}; @@ -222,13 +219,18 @@ where pub fn attach_node(&self, node: Node) { self.set.attach_node(node) } - /// Returns detached snapshots of this map's [`Node`]s. + /// Attaches persisted [`Node`]s with one topology publication. + #[cfg(feature = "cdc")] + pub fn attach_nodes(&self, nodes: impl IntoIterator) { + self.set.attach_nodes(nodes) + } + + /// Returns detached, read-only snapshots of this map's [`Node`]s. /// - /// The returned mutexes preserve the checkpoint API's existing shape, but - /// mutating them does not mutate this map. Callers requiring one coherent - /// logical generation must prevent concurrent mutation while collecting. + /// Callers requiring one coherent logical generation must prevent + /// concurrent mutation while collecting. #[cfg(feature = "cdc")] - pub fn iter_nodes(&self) -> impl Iterator>> + '_ + pub fn snapshot_nodes(&self) -> Vec where Node: Clone, { @@ -236,9 +238,8 @@ where .index .read() .values() - .map(|node| Arc::new(Mutex::new((*node.read()).clone()))) - .collect::>() - .into_iter() + .map(|node| (*node.read()).clone()) + .collect() } /// Copies the exact node boundaries into a pointer-free checkpoint image. @@ -287,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. diff --git a/src/concurrent/multimap.rs b/src/concurrent/multimap.rs index c115aeb..b166582 100644 --- a/src/concurrent/multimap.rs +++ b/src/concurrent/multimap.rs @@ -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, @@ -210,13 +207,18 @@ where pub fn attach_multi_node(&self, node: Node) { self.set.attach_node(node) } - /// Returns detached snapshots of this multimap's [`Node`]s. + /// Attaches persisted [`Node`]s with one topology publication. + #[cfg(feature = "cdc")] + pub fn attach_multi_nodes(&self, nodes: impl IntoIterator) { + self.set.attach_nodes(nodes) + } + + /// Returns detached, read-only snapshots of this multimap's [`Node`]s. /// - /// The returned mutexes preserve the checkpoint API's existing shape, but - /// mutating them does not mutate this map. Callers requiring one coherent - /// logical generation must prevent concurrent mutation while collecting. + /// Callers requiring one coherent logical generation must prevent + /// concurrent mutation while collecting. #[cfg(feature = "cdc")] - pub fn iter_nodes(&self) -> impl Iterator>> + '_ + pub fn snapshot_nodes(&self) -> Vec where Node: Clone, { @@ -224,9 +226,8 @@ where .index .read() .values() - .map(|node| Arc::new(Mutex::new((*node.read()).clone()))) - .collect::>() - .into_iter() + .map(|node| (*node.read()).clone()) + .collect() } /// Returns `true` if the map contains at least one occurance of the specified key. /// diff --git a/src/concurrent/operation.rs b/src/concurrent/operation.rs index 9c87f13..e943084 100644 --- a/src/concurrent/operation.rs +++ b/src/concurrent/operation.rs @@ -2,9 +2,9 @@ use std::fmt::Debug; use std::sync::Arc; use parking_lot::RwLock; -use std::collections::BTreeMap; use crate::cdc::change::ChangeEventUnassigned; +use crate::concurrent::set::TopologyWriteGuard; use crate::core::node::NodeLike; type OldVersion = Arc>; @@ -27,7 +27,7 @@ where // value in place: see `MultiPairLike::adopt_stored_identity`. pub fn commit( self, - index: &mut BTreeMap>>, + index: &mut TopologyWriteGuard<'_, T, Node>, adopt: fn(&T, &mut T), ) -> Result<(Option, Vec>), ()> { match self { @@ -166,6 +166,12 @@ where }; cdc.push(node_removal); } + // An UpdateMax is normally a route-preserving + // rekey and skips publication. A racing drain + // changes it into a true node removal, so opt + // this guard back into publication before the + // canonical entry is unlinked. + index.enable_publication(); index.remove(&old_max); (None, cdc) @@ -174,6 +180,14 @@ where // in either direction. Some(new_max) if *new_max != old_max => { let new_max = new_max.clone(); + // A stale boundary is definitive only for the + // last node, which is the point-read fallback. + // For any earlier node, publish its new route: + // after a shrink, another insert may fill the + // gap in the following node. + if !index.is_last_node(&node) { + index.enable_publication(); + } index.remove(&old_max); index.insert(new_max, node.clone()); diff --git a/src/concurrent/ref.rs b/src/concurrent/ref.rs index 0bf31d4..4b5c779 100644 --- a/src/concurrent/ref.rs +++ b/src/concurrent/ref.rs @@ -2,6 +2,12 @@ use crate::core::node::NodeLike; use parking_lot::{ArcRwLockReadGuard, RawRwLock}; use std::marker::PhantomData; +/// A point reference that keeps its node read-locked. +/// +/// Drop this value before calling another operation on the same map from the +/// same thread. Like every non-recursive `RwLock`, reacquiring a read lock while +/// a writer is queued can deadlock even when the thread already holds a read +/// guard. pub struct Ref + Send> { pub(super) node_guard: ArcRwLockReadGuard, pub(super) position: usize, diff --git a/src/concurrent/set.rs b/src/concurrent/set.rs index f29cfe8..c9d72cf 100644 --- a/src/concurrent/set.rs +++ b/src/concurrent/set.rs @@ -1,5 +1,7 @@ -use parking_lot::{ArcRwLockReadGuard, ArcRwLockWriteGuard, RawRwLock, RwLock, RwLockReadGuard, RwLockWriteGuard}; -use std::collections::BTreeMap; +use parking_lot::{ + ArcRwLockReadGuard, ArcRwLockWriteGuard, Mutex, MutexGuard, RawRwLock, RwLock, RwLockReadGuard, RwLockWriteGuard, +}; +use std::collections::{BTreeMap, HashMap}; use std::fmt::Debug; use std::iter::FusedIterator; use std::marker::PhantomData; @@ -15,16 +17,181 @@ use crate::core::node::*; use super::r#ref::Ref; const ROOT_PUBLICATION_SPIN_LIMIT: usize = 16; +const STABLE_READ_BLOCKING_FALLBACK_AFTER: usize = 2; +const PUBLICATION_BACKLOG_DRAIN_THRESHOLD: usize = 64; type NodeIndex = BTreeMap>>; -struct RetiredIndex(*mut NodeIndex); +// Point-read routes are kept in immutable, cache-friendly chunks. Publishing +// clones only the short vector of chunk Arcs and the one chunk containing the +// changed boundary; it never copies the full node index. At WorkTable's +// default 1,024 rows per node, one 128-route chunk covers roughly 131k rows. +const PUBLISHED_ROUTES_PER_CHUNK: usize = 128; + +struct PublishedChunk { + entries: Vec<(T, Arc>)>, +} + +impl Clone for PublishedChunk { + fn clone(&self) -> Self { + Self { + entries: self.entries.clone(), + } + } +} + +struct PublishedNodeIndex { + chunks: Vec>>, + len: usize, +} + +impl Clone for PublishedNodeIndex { + fn clone(&self) -> Self { + Self { + chunks: self.chunks.clone(), + len: self.len, + } + } +} + +impl PublishedNodeIndex +where + T: Ord + Clone, +{ + fn iter(&self) -> impl Iterator>)> { + self.chunks + .iter() + .flat_map(|chunk| chunk.entries.iter().map(|(key, node)| (key, node))) + } + + fn first_key_value(&self) -> Option<(&T, &Arc>)> { + self.chunks.first()?.entries.first().map(|(key, node)| (key, node)) + } + + fn last_key_value(&self) -> Option<(&T, &Arc>)> { + self.chunks.last()?.entries.last().map(|(key, node)| (key, node)) + } + + fn chunk_for(&self, key: &Q) -> usize + where + T: Borrow, + Q: Ord + ?Sized, + { + self.chunks.partition_point(|chunk| { + let max = &chunk.entries.last().expect("published chunks are non-empty").0; + >::borrow(max) < key + }) + } + + fn first_for_bound(&self, bound: Bound<&Q>) -> Option<(&T, &Arc>)> + where + T: Borrow, + Q: Ord + ?Sized, + { + let key = match bound { + Bound::Included(key) | Bound::Excluded(key) => key, + Bound::Unbounded => return self.first_key_value(), + }; + let mut chunk_index = self.chunk_for(key); + while let Some(chunk) = self.chunks.get(chunk_index) { + let entry_index = chunk.entries.partition_point(|(candidate, _)| match bound { + Bound::Included(_) => >::borrow(candidate) < key, + Bound::Excluded(_) => >::borrow(candidate) <= key, + Bound::Unbounded => false, + }); + if let Some((found, node)) = chunk.entries.get(entry_index) { + return Some((found, node)); + } + chunk_index += 1; + } + None + } + + fn insert(&mut self, key: T, node: Arc>) -> Option>> { + if self.chunks.is_empty() { + self.chunks.push(Arc::new(PublishedChunk { + entries: vec![(key, node)], + })); + self.len = 1; + return None; + } + + let mut chunk_index = self.chunk_for(&key); + if chunk_index == self.chunks.len() { + chunk_index -= 1; + } + let chunk = Arc::make_mut(&mut self.chunks[chunk_index]); + match chunk.entries.binary_search_by(|(candidate, _)| candidate.cmp(&key)) { + Ok(index) => Some(std::mem::replace(&mut chunk.entries[index].1, node)), + Err(index) => { + chunk.entries.insert(index, (key, node)); + self.len += 1; + if chunk.entries.len() > PUBLISHED_ROUTES_PER_CHUNK { + let right = chunk.entries.split_off(chunk.entries.len() / 2); + self.chunks + .insert(chunk_index + 1, Arc::new(PublishedChunk { entries: right })); + } + None + } + } + } + + fn remove(&mut self, key: &Q) -> Option>> + where + T: Borrow, + Q: Ord + ?Sized, + { + let chunk_index = self.chunk_for(key); + let chunk = self.chunks.get_mut(chunk_index)?; + let chunk = Arc::make_mut(chunk); + let entry_index = chunk + .entries + .binary_search_by(|(candidate, _)| >::borrow(candidate).cmp(key)) + .ok()?; + let (_, removed) = chunk.entries.remove(entry_index); + self.len -= 1; + + if chunk.entries.is_empty() { + self.chunks.remove(chunk_index); + } else if chunk_index > 0 + && self.chunks[chunk_index - 1].entries.len() + self.chunks[chunk_index].entries.len() + <= PUBLISHED_ROUTES_PER_CHUNK + { + let right = self.chunks.remove(chunk_index); + Arc::make_mut(&mut self.chunks[chunk_index - 1]) + .entries + .extend(right.entries.iter().cloned()); + } else if chunk_index + 1 < self.chunks.len() + && self.chunks[chunk_index].entries.len() + self.chunks[chunk_index + 1].entries.len() + <= PUBLISHED_ROUTES_PER_CHUNK + { + let right = self.chunks.remove(chunk_index + 1); + Arc::make_mut(&mut self.chunks[chunk_index]) + .entries + .extend(right.entries.iter().cloned()); + } + + Some(removed) + } +} + +#[inline] +fn node_identity(node: &Arc>) -> usize { + // Identity token only: it is never converted back into or dereferenced as + // a pointer. The Arc stays live while the token is present, preventing + // allocator reuse from aliasing two published nodes. + Arc::as_ptr(node) as usize +} + +struct RetiredIndex(*mut PublishedNodeIndex); // SAFETY: the pointer is uniquely owned after it has been swapped out of the // publication slot, and this wrapper exposes no access to the map. Its only -// operation is destruction after the grace period. Dropping the keys and, for -// the final Arc, moving each node into its destructor on that thread is valid -// when both stored types are `Send`; sharing `Node` there is not required. +// operation is destruction after the grace period. Dropping a shared route +// chunk only decrements its Arc; dropping the final route path can move/drop T +// and Node on the reclaiming thread, hence Send. The wrapper never dereferences +// the index, and its private field prevents callers from adding such access +// without revisiting this proof. unsafe impl Send for RetiredIndex {} impl Drop for RetiredIndex { @@ -35,16 +202,18 @@ impl Drop for RetiredIndex { } } -#[derive(Debug)] struct PublishedIndex { - current: AtomicPtr>, + current: AtomicPtr>, domain: ps_reclaim::Domain, } impl PublishedIndex { fn new() -> Self { Self { - current: AtomicPtr::new(Box::into_raw(Box::new(BTreeMap::new()))), + current: AtomicPtr::new(Box::into_raw(Box::new(PublishedNodeIndex { + chunks: Vec::new(), + len: 0, + }))), domain: ps_reclaim::Domain::new(), } } @@ -55,16 +224,37 @@ where T: Ord + Clone + Send + 'static, Node: Send + 'static, { - fn publish(&self, index: &NodeIndex) { - let replacement = Box::into_raw(Box::new(index.clone())); + fn snapshot(&self) -> PublishedNodeIndex { + let current = self.current.load(Ordering::Acquire); + // SAFETY: callers hold the only structural writer lock. `current` + // cannot be unlinked until that writer publishes its replacement. + unsafe { (&*current).clone() } + } + + fn replace(&self, replacement: PublishedNodeIndex) -> RetiredIndex { + // The route index is structurally shared: publishing moves one root, + // and the writer copied only its chunk-Arc vector plus touched chunks. + let replacement = Box::into_raw(Box::new(replacement)); let retired = self.current.swap(replacement, Ordering::AcqRel); + RetiredIndex(retired) + } + + fn retire(&self, retired: RetiredIndex) { // Keep the pointer's provenance intact while transferring its unique // ownership to the retirement callback. - let retired = RetiredIndex(retired); self.domain.retire(move || drop(retired)); - // Topology publication is infrequent (normally one node boundary per - // 1,024 inserts), so reclaim one expired snapshot on the writer path. - self.domain.advance_up_to(1); + } + + fn advance(&self) { + // A reader delayed on a node writer never holds a pin (see the point + // read paths below). Do not sweep ps-reclaim's 256-slot registry on + // every split: that turns the registry into the same reader/writer + // cache-line fight this publication path removes. A small bounded + // backlog amortizes the sweep while `advance` drains every route root + // whose grace period has elapsed. + if self.domain.pending() >= PUBLICATION_BACKLOG_DRAIN_THRESHOLD { + self.domain.advance(); + } } } @@ -77,19 +267,30 @@ impl Drop for PublishedIndex { } } -#[derive(Debug)] pub(crate) struct Topology { index: RwLock>, + published_keys: Mutex>, published: PublishedIndex, // Even values are stable publications; odd values mean a writer may have // changed node contents or routing but has not published the new route. generation: AtomicU64, } +impl Debug for Topology { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("Topology") + .field("nodes", &self.index.read().len()) + .field("generation", &self.generation.load(Ordering::Relaxed)) + .finish() + } +} + impl Topology { fn new() -> Self { Self { index: RwLock::new(BTreeMap::new()), + published_keys: Mutex::new(HashMap::new()), published: PublishedIndex::new(), generation: AtomicU64::new(0), } @@ -110,24 +311,64 @@ where fn write(&self) -> TopologyWriteGuard<'_, T, Node> { let index = self.index.write(); self.generation.fetch_add(1, Ordering::AcqRel); - TopologyWriteGuard { topology: self, index } + TopologyWriteGuard { + topology: self, + index: Some(index), + published: None, + published_keys: None, + publish: true, + dirty: false, + } + } + + /// Re-keys a node without replacing the point-read snapshot. + /// + /// This guard starts without a replacement snapshot. The commit may retain + /// that fast path only when re-keying the last node: point reads already + /// fall back to the last route beyond its published maximum. Re-keying any + /// earlier node must call `enable_publication`, because a subsequent insert + /// can fill a gap left by a shrinking maximum in the following node. + #[inline] + fn write_rekey(&self) -> TopologyWriteGuard<'_, T, Node> { + let index = self.index.write(); + TopologyWriteGuard { + topology: self, + index: Some(index), + published: None, + published_keys: None, + publish: false, + dirty: false, + } } #[inline] fn try_write(&self) -> Option> { let index = self.index.try_write()?; self.generation.fetch_add(1, Ordering::AcqRel); - Some(TopologyWriteGuard { topology: self, index }) + Some(TopologyWriteGuard { + topology: self, + index: Some(index), + published: None, + published_keys: None, + publish: true, + dirty: false, + }) } } -struct TopologyWriteGuard<'a, T, Node> +pub(crate) struct TopologyWriteGuard<'a, T, Node> where T: Ord + Clone + Send + 'static, Node: Send + 'static, { topology: &'a Topology, - index: RwLockWriteGuard<'a, NodeIndex>, + // Option lets Drop release the structural lock before advancing the + // reclamation domain. Range readers should not wait for a registry scan. + index: Option>>, + published: Option>, + published_keys: Option>>, + publish: bool, + dirty: bool, } impl std::ops::Deref for TopologyWriteGuard<'_, T, Node> @@ -138,17 +379,130 @@ where type Target = NodeIndex; fn deref(&self) -> &Self::Target { - &self.index + self.index.as_deref().expect("topology guard already released") } } -impl std::ops::DerefMut for TopologyWriteGuard<'_, T, Node> +impl<'a, T, Node> TopologyWriteGuard<'a, T, Node> where T: Ord + Clone + Send + 'static, Node: Send + 'static, { - fn deref_mut(&mut self) -> &mut Self::Target { - &mut self.index + pub(crate) fn enable_publication(&mut self) { + if self.publish { + return; + } + // `write_rekey` deliberately leaves the stable generation untouched + // for a route-safe last-node rekey. If commit discovers that the + // route really must change, enter the odd writer generation before + // constructing or publishing its replacement. + self.topology.generation.fetch_add(1, Ordering::AcqRel); + self.publish = true; + } + + fn ensure_publication_snapshot(&mut self) { + if self.published.is_none() { + self.published_keys = Some(self.topology.published_keys.lock()); + self.published = Some(self.topology.published.snapshot()); + } + } + + pub(crate) fn is_last_node(&self, node: &Arc>) -> bool { + self.index + .as_deref() + .and_then(BTreeMap::last_key_value) + .is_some_and(|(_, candidate)| Arc::ptr_eq(candidate, node)) + } + + pub(crate) fn insert(&mut self, key: T, node: Arc>) -> Option>> { + let replaced = self + .index + .as_deref_mut() + .expect("topology guard already released") + .insert(key.clone(), node.clone()); + self.dirty = true; + + if self.publish { + self.ensure_publication_snapshot(); + let published = self.published.as_mut().expect("publication snapshot initialized"); + let published_keys = self + .published_keys + .as_mut() + .expect("publication identity map initialized"); + if let Some(replaced) = &replaced { + if let Some(old_key) = published_keys.remove(&node_identity(replaced)) { + published.remove(&old_key); + } + } + + // A skipped rekey can leave another node published at the new + // canonical boundary. Move that displaced node to its own current + // canonical key, repeating only if stale boundaries form a short + // collision chain. The fallback scan is confined to such a + // collision; ordinary split/attach publication remains O(log N). + let canonical = self.index.as_deref().expect("topology guard already released"); + let mut route_key = key; + let mut route_node = node; + loop { + let displaced = published.insert(route_key.clone(), route_node.clone()); + published_keys.insert(node_identity(&route_node), route_key.clone()); + + let Some(displaced) = displaced else { + break; + }; + if Arc::ptr_eq(&displaced, &route_node) { + break; + } + published_keys.remove(&node_identity(&displaced)); + + let Some((canonical_key, _)) = canonical + .iter() + .find(|(_, candidate)| Arc::ptr_eq(candidate, &displaced)) + else { + // The displaced route belonged to a node removed from the + // canonical topology by an earlier route-preserving race. + break; + }; + route_key = canonical_key.clone(); + route_node = displaced; + } + } + + replaced + } + + pub(crate) fn remove(&mut self, key: &Q) -> Option>> + where + T: Borrow, + Q: Ord + ?Sized, + { + let removed = self + .index + .as_deref_mut() + .expect("topology guard already released") + .remove(key)?; + self.dirty = true; + + if self.publish { + self.ensure_publication_snapshot(); + let published = self.published.as_mut().expect("publication snapshot initialized"); + let published_keys = self + .published_keys + .as_mut() + .expect("publication identity map initialized"); + let route_key = published_keys + .remove(&node_identity(&removed)) + .expect("published key map must contain a canonical node"); + let route_node = published + .remove::(&route_key) + .expect("published route must contain its recorded boundary"); + assert!( + Arc::ptr_eq(&route_node, &removed), + "published boundary must identify the removed node" + ); + } + + Some(removed) } } @@ -158,8 +512,56 @@ where Node: Send + 'static, { fn drop(&mut self) { - self.topology.published.publish(&self.index); + #[cfg(debug_assertions)] + if self.publish && self.dirty { + let canonical = self.index.as_deref().expect("topology guard exists while validating"); + let published = self.published.as_ref().expect("publishing guard carries snapshot"); + let published_keys = self + .published_keys + .as_deref() + .expect("publishing guard carries identity map"); + debug_assert_eq!( + canonical.len(), + published.len, + "canonical and published node counts diverged" + ); + debug_assert_eq!( + canonical.len(), + published_keys.len(), + "canonical and published identity counts diverged" + ); + for node in canonical.values() { + debug_assert!( + published_keys.contains_key(&node_identity(node)), + "canonical node is absent from published identity map" + ); + } + } + + if !self.publish { + drop(self.published_keys.take()); + drop(self.index.take()); + return; + } + if !self.dirty { + self.topology.generation.fetch_add(1, Ordering::Release); + drop(self.published_keys.take()); + drop(self.index.take()); + return; + } + let retired = self.topology.published.replace( + self.published + .take() + .expect("publishing guard must carry a read snapshot"), + ); + // Readers may proceed as soon as the O(1) root publication completes. + // Garbage bookkeeping and the registry sweep are deliberately outside + // that odd-generation window. self.topology.generation.fetch_add(1, Ordering::Release); + self.topology.published.retire(retired); + drop(self.published_keys.take()); + drop(self.index.take()); + self.topology.published.advance(); } } @@ -194,6 +596,26 @@ where }) } +fn first_published_for_borrowed_bound<'a, T, Q, V>( + index: &'a PublishedNodeIndex, + bound: Bound<&Q>, + borrow_order_matches: bool, +) -> Option<(&'a T, &'a Arc>)> +where + T: Ord + Clone + Borrow, + Q: Ord + ?Sized, +{ + if borrow_order_matches { + return index.first_for_bound(bound); + } + + index.iter().find(|(key, _)| match bound { + Bound::Included(value) => >::borrow(key) >= value, + Bound::Excluded(value) => >::borrow(key) > value, + Bound::Unbounded => true, + }) +} + fn node_for_borrowed_end<'a, T, Q, V>( index: &'a BTreeMap, end: &Q, @@ -345,11 +767,27 @@ where self } pub fn attach_node(&self, node: Node) { - let node_id = node - .max() - .cloned() - .expect("node should contain at least one value to be correct node"); - self.index.write().insert(node_id, Arc::new(RwLock::new(node))); + self.attach_nodes(std::iter::once(node)); + } + + /// Attaches a persisted topology in one structural publication. + /// + /// Nodes must be non-empty, internally sorted, and mutually ordered. The + /// same preconditions as [`Self::attach_node`] apply to every item. + pub fn attach_nodes(&self, nodes: impl IntoIterator) { + let mut nodes = nodes.into_iter().peekable(); + if nodes.peek().is_none() { + return; + } + + let mut index = self.index.write(); + for node in nodes { + let node_id = node + .max() + .cloned() + .expect("node should contain at least one value to be correct node"); + index.insert(node_id, Arc::new(RwLock::new(node))); + } } #[cfg(feature = "cdc")] @@ -475,9 +913,11 @@ where drop(node_guard); drop(index); - let mut index = self.index.write(); - let op = operation.unwrap(); + let mut index = match &op { + Operation::UpdateMax(_, _) => self.index.write_rekey(), + Operation::Split(_, _, _) | Operation::MakeUnreachable(_, _) => self.index.write(), + }; match &op { Operation::Split(_, _, _) => { if let Ok((value, value_cdc)) = op.commit::(&mut index, adopt) { @@ -708,9 +1148,13 @@ where drop(node_guard); drop(index); - let mut index = self.index.write(); + let operation = operation.unwrap(); + let mut index = match &operation { + Operation::UpdateMax(_, _) => self.index.write_rekey(), + Operation::Split(_, _, _) | Operation::MakeUnreachable(_, _) => self.index.write(), + }; - return if let Ok((_, value_cdc)) = operation.unwrap().commit::(&mut index, no_identity_adoption) { + return if let Ok((_, value_cdc)) = operation.commit::(&mut index, no_identity_adoption) { #[cfg(feature = "cdc")] if EMIT_CDC { for unassigned_event in value_cdc { @@ -796,37 +1240,69 @@ where T: Borrow, Q: Ord + ?Sized, { + let mut retries = 0; + let mut writer_spins = 0; + loop { let generation = self.index.generation.load(Ordering::Acquire); if !generation.is_multiple_of(2) { - std::hint::spin_loop(); + if writer_spins < ROOT_PUBLICATION_SPIN_LIMIT { + writer_spins += 1; + std::hint::spin_loop(); + } else { + std::thread::yield_now(); + } continue; } + writer_spins = 0; + + if retries >= STABLE_READ_BLOCKING_FALLBACK_AFTER { + // Bounded progress fallback: hold the canonical topology read + // guard while acquiring the node. Structural writers follow + // the same topology-before-node order. + let index = self.index.read(); + let node = match first_for_borrowed_bound(&index, Bound::Included(value), self.borrow_order_matches) { + Some((_, node)) => Some(node.clone()), + None => index + .last_key_value() + .map(|(_, node)| node.clone()) + .or_else(|| index.first_key_value().map(|(_, node)| node.clone())), + }?; + let node_guard = node.read_arc(); + drop(index); + return Some(node_guard); + } let pin = self.index.published.domain.pin(); let snapshot = self.index.published.current.load(Ordering::Acquire); // SAFETY: `snapshot` was loaded after `pin`, and the publication // domain cannot reclaim it until `pin` is dropped. let index = unsafe { &*snapshot }; - let node = match first_for_borrowed_bound(index, Bound::Included(value), self.borrow_order_matches) { - Some((_, node)) => Some(node.clone()), - None => index - .last_key_value() - .map(|(_, node)| node.clone()) - .or_else(|| index.first_key_value().map(|(_, node)| node.clone())), - }; + let node = + match first_published_for_borrowed_bound(index, Bound::Included(value), self.borrow_order_matches) { + Some((_, node)) => Some(node.clone()), + None => index + .last_key_value() + .map(|(_, node)| node.clone()) + .or_else(|| index.first_key_value().map(|(_, node)| node.clone())), + }; let Some(node) = node else { if self.index.generation.load(Ordering::Acquire) == generation { return None; } + retries += 1; continue; }; + // The snapshot pin protects the Arc only until it is cloned. Drop + // it before the potentially blocking node acquisition so a slow + // node writer cannot stall topology reclamation. + drop(pin); let node_guard = node.read_arc(); if self.index.generation.load(Ordering::Acquire) == generation { - drop(pin); return Some(node_guard); } + retries += 1; } } @@ -850,19 +1326,42 @@ where T: Borrow, Q: Ord + ?Sized, { + let mut retries = 0; + let mut writer_spins = 0; + let mut read = Some(read); + loop { let generation = self.index.generation.load(Ordering::Acquire); if !generation.is_multiple_of(2) { - std::hint::spin_loop(); + if writer_spins < ROOT_PUBLICATION_SPIN_LIMIT { + writer_spins += 1; + std::hint::spin_loop(); + } else { + std::thread::yield_now(); + } continue; } + writer_spins = 0; + + if retries >= STABLE_READ_BLOCKING_FALLBACK_AFTER { + let index = self.index.read(); + let node = first_for_borrowed_bound(&index, Bound::Included(value), self.borrow_order_matches) + .or_else(|| index.last_key_value()) + .or_else(|| index.first_key_value()) + .map(|(_, node)| node)?; + let node_guard = node.read(); + let position = node_guard.try_select(value)?; + return node_guard + .get_ith(position) + .map(read.take().expect("read closure is consumed only on return")); + } - let _pin = self.index.published.domain.pin(); + let pin = self.index.published.domain.pin(); let snapshot = self.index.published.current.load(Ordering::Acquire); // SAFETY: `snapshot` was loaded after `pin`, and the publication // domain cannot reclaim it until `pin` is dropped. let index = unsafe { &*snapshot }; - let node = first_for_borrowed_bound(index, Bound::Included(value), self.borrow_order_matches) + let node = first_published_for_borrowed_bound(index, Bound::Included(value), self.borrow_order_matches) .or_else(|| index.last_key_value()) .or_else(|| index.first_key_value()) .map(|(_, node)| node); @@ -870,18 +1369,37 @@ where if self.index.generation.load(Ordering::Acquire) == generation { return None; } + retries += 1; continue; }; // Borrow the Arc from the protected snapshot: unlike // `lock_node_for_value`, this owned-result path does not need an - // Arc clone or its shared refcount update. - let node_guard = node.read(); + // Arc clone or its shared refcount update when the node is free. + // On contention, clone it and release the reclamation pin before + // blocking so a node writer cannot retain old topology paths. + if let Some(node_guard) = node.try_read() { + if self.index.generation.load(Ordering::Acquire) != generation { + retries += 1; + continue; + } + let position = node_guard.try_select(value)?; + return node_guard + .get_ith(position) + .map(read.take().expect("read closure is consumed only on return")); + } + + let node = node.clone(); + drop(pin); + let node_guard = node.read_arc(); if self.index.generation.load(Ordering::Acquire) != generation { + retries += 1; continue; } let position = node_guard.try_select(value)?; - return node_guard.get_ith(position).map(read); + return node_guard + .get_ith(position) + .map(read.take().expect("read closure is consumed only on return")); } } @@ -1789,7 +2307,7 @@ mod tests { #[test] fn published_point_reads_remain_definitive_across_splits() { - use std::sync::atomic::{AtomicBool, Ordering}; + use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; const STABLE_KEYS: usize = 256; const FINAL_KEYS: usize = 4_096; @@ -1802,17 +2320,25 @@ mod tests { let start = Arc::new(Barrier::new(READERS + 1)); let done = Arc::new(AtomicBool::new(false)); + let published_up_to = Arc::new(AtomicUsize::new(STABLE_KEYS - 1)); let readers = (0..READERS) .map(|reader| { let set = Arc::clone(&set); let start = Arc::clone(&start); let done = Arc::clone(&done); + let published_up_to = Arc::clone(&published_up_to); thread::spawn(move || { start.wait(); let mut probe = reader; while !done.load(Ordering::Acquire) { let key = probe % STABLE_KEYS; assert_eq!(set.get_with(&key, |value| *value), Some(key)); + let newest_acknowledged = published_up_to.load(Ordering::Acquire); + assert_eq!( + set.get_with(&newest_acknowledged, |value| *value), + Some(newest_acknowledged), + "an acknowledged insert disappeared from the published route" + ); probe += READERS; } }) @@ -1822,6 +2348,7 @@ mod tests { start.wait(); for key in STABLE_KEYS..FINAL_KEYS { set.insert(key); + published_up_to.store(key, Ordering::Release); } done.store(true, Ordering::Release); @@ -1829,6 +2356,66 @@ mod tests { reader.join().unwrap(); } assert_eq!(set.len(), FINAL_KEYS); + for key in 0..FINAL_KEYS { + assert_eq!(set.get_with(&key, |value| *value), Some(key)); + } + } + + #[test] + fn published_pointer_survives_reclamation_interleavings() { + use std::sync::atomic::{AtomicBool, Ordering}; + + let set = Arc::new(BTreeSet::::with_maximum_node_size(2)); + for key in 0..8 { + set.insert(key); + } + let start = Arc::new(Barrier::new(2)); + let done = Arc::new(AtomicBool::new(false)); + + let reader_set = Arc::clone(&set); + let reader_start = Arc::clone(&start); + let reader_done = Arc::clone(&done); + let reader = thread::spawn(move || { + reader_start.wait(); + let mut probe = 0; + while !reader_done.load(Ordering::Acquire) { + let key = probe % 8; + assert_eq!(reader_set.get_with(&key, |value| *value), Some(key)); + probe += 1; + } + }); + + start.wait(); + for key in 8..80 { + set.insert(key); + } + for key in 8..80 { + assert_eq!(set.remove(&key), Some(key)); + } + done.store(true, Ordering::Release); + reader.join().unwrap(); + + for key in 0..8 { + assert_eq!(set.get_with(&key, |value| *value), Some(key)); + } + } + + #[test] + fn published_route_chunks_split_and_merge_without_losing_keys() { + let set = BTreeSet::::with_maximum_node_size(2); + for key in 0..600 { + assert!(set.insert(key)); + } + for key in (0..600).step_by(2) { + assert_eq!(set.remove(&key), Some(key)); + } + for key in 0..600 { + assert_eq!(set.contains(&key), key % 2 == 1, "probe {key}"); + } + for key in (1..600).step_by(2) { + assert_eq!(set.remove(&key), Some(key)); + } + assert!(set.is_empty()); } #[test] @@ -2460,6 +3047,58 @@ mod tests { } } + #[test] + fn published_route_updates_when_a_non_last_boundary_shrinks() { + let set = BTreeSet::::with_maximum_node_size(8); + set.attach_nodes([vec![1, 10], vec![20, 30]]); + + // Removing the first node's maximum changes its canonical route from + // 10 to 1. If the published route remains at 10, the later insertion + // of 5 correctly lands in the second node but a point read for 5 is + // misrouted to the first node and reports a false miss. + assert_eq!(set.remove(&10), Some(10)); + assert!(set.insert(5)); + assert!(set.contains(&5)); + assert_eq!(set.get(&5).map(|value| *value.get()), Some(5)); + } + + #[test] + fn published_point_routes_match_a_sequential_oracle_under_churn() { + let set = BTreeSet::::with_maximum_node_size(4); + let mut oracle = std::collections::BTreeSet::new(); + let mut state = 0x8f4d_2a71_c390_6be5u64; + + for step in 0..2_000 { + // Fixed xorshift stream: deterministic inserts/removes repeatedly + // grow, shrink, empty, and split tiny nodes. + state ^= state << 13; + state ^= state >> 7; + state ^= state << 17; + let key = state % 64; + if state & 1 == 0 { + assert_eq!(set.insert(key), oracle.insert(key), "insert step {step}, key {key}"); + } else { + assert_eq!( + set.remove(&key).is_some(), + oracle.remove(&key), + "remove step {step}, key {key}" + ); + } + + for probe in 0..64 { + assert_eq!( + set.contains(&probe), + oracle.contains(&probe), + "point route diverged at step {step}, probe {probe}" + ); + } + assert_eq!( + set.iter().collect::>(), + oracle.iter().copied().collect::>() + ); + } + } + #[test] fn concurrent_remove_reinsert_over_emptying_nodes_preserves_all_keys() { const THREADS: u64 = 4; diff --git a/src/core/node.rs b/src/core/node.rs index 4d3dc9b..5fcd5e9 100644 --- a/src/core/node.rs +++ b/src/core/node.rs @@ -25,6 +25,8 @@ pub trait NodeLike { where T: Borrow; #[allow(dead_code)] + /// Must return `Some(i)` exactly when [`Self::contains`] is true, with + /// `get_ith(i)` equal to the requested value. fn try_select(&self, value: &Q) -> Option where T: Borrow; @@ -70,11 +72,6 @@ mod search_backend { { haystack.binary_search_by(|candidate| candidate.borrow().cmp(needle)) } - - #[inline] - pub(crate) fn search_by(haystack: &[T], compare: impl FnMut(&T) -> core::cmp::Ordering) -> Result { - haystack.binary_search_by(compare) - } } #[cfg(all( @@ -104,18 +101,6 @@ mod search_backend { _ => Err(index), } } - - #[inline] - pub(crate) fn search_by( - haystack: &[T], - mut compare: impl FnMut(&T) -> core::cmp::Ordering, - ) -> Result { - let index = haystack.lower_bound_by(&mut compare); - match haystack.get(index) { - Some(candidate) if compare(candidate).is_eq() => Ok(index), - _ => Err(index), - } - } } #[cfg(all( @@ -137,11 +122,6 @@ mod search_backend { { haystack.exact_search_by(|candidate| candidate.borrow().cmp(needle)) } - - #[inline] - pub(crate) fn search_by(haystack: &[T], compare: impl FnMut(&T) -> core::cmp::Ordering) -> Result { - haystack.exact_search_by(compare) - } } #[cfg(any( @@ -189,26 +169,6 @@ mod search_backend { Err(i) } - - #[inline] - pub(crate) fn search_by(haystack: &[T], mut compare: impl FnMut(&T) -> Ordering) -> Result { - let mut right = haystack.len(); - let mut left = 0; - - while left != right { - let middle = (left + right) >> 1; - // SAFETY: `left < right <= haystack.len()` makes `middle` a valid - // element index, and both branches preserve the bounds. - let candidate = unsafe { haystack.get_unchecked(middle) }; - match compare(candidate) { - Ordering::Equal => return Ok(middle), - Ordering::Less => left = middle + 1, - Ordering::Greater => right = middle, - } - } - - Err(left) - } } // Search backend precedence is deterministic when features are composed: @@ -216,7 +176,20 @@ mod search_backend { // custom implementation is the compatibility fallback. Each backend's cfg // selects its implementation and test name together, preventing drift. use search_backend::search; -pub(crate) use search_backend::search_by; + +/// Returns the first comparator-equal entry, or its insertion position. +/// +/// This helper deliberately has one implementation across configured search +/// backends: callers that compare only a prefix (such as map key without +/// value) must not observe backend-dependent positions among duplicates. +#[inline] +pub(crate) fn search_by(haystack: &[T], mut compare: impl FnMut(&T) -> core::cmp::Ordering) -> Result { + let index = haystack.partition_point(|candidate| compare(candidate).is_lt()); + match haystack.get(index) { + Some(candidate) if compare(candidate).is_eq() => Ok(index), + _ => Err(index), + } +} #[inline] fn compute_positions_to_skip(haystack: &[T], bound: std::ops::Bound<&Q>, forward: bool) -> Option @@ -397,6 +370,14 @@ mod tests { assert_eq!(search_backend::NAME, "custom"); } + #[test] + fn comparator_search_returns_first_duplicate() { + let values = [(1, "a"), (1, "b"), (1, "c"), (2, "d")]; + assert_eq!(search_by(&values, |candidate| candidate.0.cmp(&1)), Ok(0)); + assert_eq!(search_by(&values, |candidate| candidate.0.cmp(&2)), Ok(3)); + assert_eq!(search_by(&values, |candidate| candidate.0.cmp(&0)), Err(0)); + } + #[test] fn point_search_returns_match_or_insertion_position() { for values in [vec![], vec![2], vec![2, 4, 8, 16, 32]] { diff --git a/src/lib.rs b/src/lib.rs index 0ef3a0c..ceb72c0 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -170,7 +170,7 @@ impl BTreeSet { // never return less than 0, so it is only necessary to check whether it is out of bounds // from the right if self.inner.get(node_idx).is_none() { - node_idx -= 1 + node_idx = node_idx.saturating_sub(1) } node_idx @@ -190,7 +190,7 @@ impl BTreeSet { }); if self.inner.get(node_idx).is_none() { - node_idx -= 1 + node_idx = node_idx.saturating_sub(1) } node_idx diff --git a/tests/loom_publication.rs b/tests/loom_publication.rs new file mode 100644 index 0000000..e2c7330 --- /dev/null +++ b/tests/loom_publication.rs @@ -0,0 +1,63 @@ +use loom::sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, +}; + +struct Publication { + generation: AtomicUsize, + current: AtomicUsize, + snapshots: [AtomicUsize; 2], +} + +/// Models the ordering contract used by `concurrent::set::Topology`. +/// +/// The writer marks the generation odd, initializes a replacement, publishes +/// it, and finally releases an even generation. A reader may accept a route +/// only when the same even generation brackets its publication load. Seeing a +/// newer publication with an older generation is harmless; seeing an older or +/// uninitialized publication after acquiring a newer generation is not. +#[test] +fn stable_even_generation_never_accepts_an_older_publication() { + loom::model(|| { + let publication = Arc::new(Publication { + generation: AtomicUsize::new(0), + current: AtomicUsize::new(0), + snapshots: [AtomicUsize::new(0), AtomicUsize::new(usize::MAX)], + }); + + let writer_publication = Arc::clone(&publication); + let writer = loom::thread::spawn(move || { + writer_publication.generation.fetch_add(1, Ordering::AcqRel); + writer_publication.snapshots[1].store(2, Ordering::Relaxed); + writer_publication.current.swap(1, Ordering::AcqRel); + writer_publication.generation.fetch_add(1, Ordering::Release); + }); + + let reader_publication = Arc::clone(&publication); + let reader = loom::thread::spawn(move || { + let before = reader_publication.generation.load(Ordering::Acquire); + if !before.is_multiple_of(2) { + return; + } + + let route = reader_publication.current.load(Ordering::Acquire); + let snapshot_generation = reader_publication.snapshots[route].load(Ordering::Relaxed); + let after = reader_publication.generation.load(Ordering::Acquire); + + if before == after { + assert_ne!( + snapshot_generation, + usize::MAX, + "an accepted publication must be fully initialized" + ); + assert!( + snapshot_generation >= before, + "an accepted publication must not predate its stable generation" + ); + } + }); + + writer.join().unwrap(); + reader.join().unwrap(); + }); +} From 2075999505f16584db1dec3ac22f6340c01104ed Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 5 Sep 2026 20:47:28 +0700 Subject: [PATCH 3/4] Harden topology publication invariants --- CHANGELOG.md | 8 + src/concurrent/operation.rs | 8 +- src/concurrent/set.rs | 323 +++++++++++++++++++++++++++++------- 3 files changed, 279 insertions(+), 60 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2054f46..464fd11 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -28,6 +28,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - 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. + ## [0.0.11] ### Changed diff --git a/src/concurrent/operation.rs b/src/concurrent/operation.rs index e943084..83ea784 100644 --- a/src/concurrent/operation.rs +++ b/src/concurrent/operation.rs @@ -185,11 +185,13 @@ where // For any earlier node, publish its new route: // after a shrink, another insert may fill the // gap in the following node. - if !index.is_last_node(&node) { + if index.is_last_node(&node) { + index.rekey_last_node(&old_max, new_max, node.clone()); + } else { index.enable_publication(); + index.remove(&old_max); + index.insert(new_max, node.clone()); } - index.remove(&old_max); - index.insert(new_max, node.clone()); (None, cdc) } diff --git a/src/concurrent/set.rs b/src/concurrent/set.rs index c9d72cf..c7d1630 100644 --- a/src/concurrent/set.rs +++ b/src/concurrent/set.rs @@ -58,6 +58,27 @@ impl PublishedNodeIndex where T: Ord + Clone, { + fn from_canonical(index: &NodeIndex) -> Self { + let mut chunks = Vec::with_capacity(index.len().div_ceil(PUBLISHED_ROUTES_PER_CHUNK)); + let mut entries = Vec::with_capacity(PUBLISHED_ROUTES_PER_CHUNK); + + for (key, node) in index { + entries.push((key.clone(), node.clone())); + if entries.len() == PUBLISHED_ROUTES_PER_CHUNK { + chunks.push(Arc::new(PublishedChunk { entries })); + entries = Vec::with_capacity(PUBLISHED_ROUTES_PER_CHUNK); + } + } + if !entries.is_empty() { + chunks.push(Arc::new(PublishedChunk { entries })); + } + + Self { + chunks, + len: index.len(), + } + } + fn iter(&self) -> impl Iterator>)> { self.chunks .iter() @@ -269,6 +290,9 @@ impl Drop for PublishedIndex { pub(crate) struct Topology { index: RwLock>, + // Writer-only reverse lookup from node identity to its current published + // route key. This differs from the canonical key only for the last node, + // whose stale route remains a valid final point-read fallback. published_keys: Mutex>, published: PublishedIndex, // Even values are stable publications; odd values mean a writer may have @@ -392,6 +416,10 @@ where if self.publish { return; } + debug_assert!( + !self.dirty, + "publication must be enabled before mutating an opt-out topology guard" + ); // `write_rekey` deliberately leaves the stable generation untouched // for a route-safe last-node rekey. If commit discovers that the // route really must change, enter the odd writer generation before @@ -407,6 +435,30 @@ where } } + /// Restores all derived publication state from the canonical topology. + /// + /// Ordinary mutations update one route in O(log N). This bounded O(N) + /// recovery is reserved for an internal identity mismatch or impossible + /// key collision; it prevents a bookkeeping defect from panicking or + /// entering an unbounded repair loop while the structural lock is held. + fn rebuild_publication(&mut self) { + if self.published_keys.is_none() { + self.published_keys = Some(self.topology.published_keys.lock()); + } + let canonical = self.index.as_deref().expect("topology guard already released"); + let rebuilt = PublishedNodeIndex::from_canonical(canonical); + let rebuilt_keys = canonical + .iter() + .map(|(key, node)| (node_identity(node), key.clone())) + .collect(); + **self + .published_keys + .as_mut() + .expect("publication identity lock was initialized") = rebuilt_keys; + self.published = Some(rebuilt); + self.dirty = true; + } + pub(crate) fn is_last_node(&self, node: &Arc>) -> bool { self.index .as_deref() @@ -414,7 +466,88 @@ where .is_some_and(|(_, candidate)| Arc::ptr_eq(candidate, node)) } + /// Re-keys the one route that may safely remain stale for point reads. + /// + /// This is deliberately separate from `insert`/`remove`: those general + /// mutation methods require publication to be enabled. The caller has + /// already verified that `old_key` identifies `node` and that it is the + /// canonical last node. Its published route remains a valid final + /// fallback whether the maximum grows or shrinks. + pub(crate) fn rekey_last_node(&mut self, old_key: &T, new_key: T, node: Arc>) { + debug_assert!(!self.publish, "last-node rekey must use the opt-out guard"); + debug_assert!( + !self.dirty, + "an opt-out guard may perform only one explicit last-node rekey" + ); + debug_assert!( + self.is_last_node(&node), + "only the canonical last node may skip publication" + ); + + let index = self.index.as_deref_mut().expect("topology guard already released"); + let removed = index.remove(old_key); + debug_assert!( + removed.as_ref().is_some_and(|removed| Arc::ptr_eq(removed, &node)), + "last-node rekey must remove its expected canonical route" + ); + let replaced = index.insert(new_key, node); + debug_assert!( + replaced.is_none(), + "last-node rekey must not collide with another canonical route" + ); + self.dirty = true; + } + + /// Makes the current canonical last-node route exact before attachment + /// can place another node after it. This runs once per attach batch, not + /// once per node. + fn repair_last_route_before_attach(&mut self) { + debug_assert!(self.publish, "attachment repair requires publication"); + let Some((canonical_key, last_node)) = self + .index + .as_deref() + .expect("topology guard already released") + .last_key_value() + .map(|(key, node)| (key.clone(), node.clone())) + else { + return; + }; + + self.ensure_publication_snapshot(); + let published_key = self + .published_keys + .as_ref() + .expect("publication identity map initialized") + .get(&node_identity(&last_node)) + .cloned(); + let Some(published_key) = published_key else { + self.rebuild_publication(); + return; + }; + if published_key == canonical_key { + return; + } + + let repaired_consistently = { + let published = self.published.as_mut().expect("publication snapshot initialized"); + let published_keys = self + .published_keys + .as_mut() + .expect("publication identity map initialized"); + let removed = published.remove(&published_key); + let displaced = published.insert(canonical_key.clone(), last_node.clone()); + published_keys.insert(node_identity(&last_node), canonical_key); + removed.is_some_and(|old_node| Arc::ptr_eq(&old_node, &last_node)) && displaced.is_none() + }; + if !repaired_consistently { + self.rebuild_publication(); + } else { + self.dirty = true; + } + } + pub(crate) fn insert(&mut self, key: T, node: Arc>) -> Option>> { + debug_assert!(self.publish, "generic topology insertion requires publication"); let replaced = self .index .as_deref_mut() @@ -424,47 +557,39 @@ where if self.publish { self.ensure_publication_snapshot(); - let published = self.published.as_mut().expect("publication snapshot initialized"); - let published_keys = self - .published_keys - .as_mut() - .expect("publication identity map initialized"); if let Some(replaced) = &replaced { - if let Some(old_key) = published_keys.remove(&node_identity(replaced)) { - published.remove(&old_key); + let removed_consistently = { + let published = self.published.as_mut().expect("publication snapshot initialized"); + let published_keys = self + .published_keys + .as_mut() + .expect("publication identity map initialized"); + published_keys + .remove(&node_identity(replaced)) + .and_then(|old_key| published.remove(&old_key)) + .is_some_and(|old_node| Arc::ptr_eq(&old_node, replaced)) + }; + if !removed_consistently { + self.rebuild_publication(); + return Some(replaced.clone()); } } - // A skipped rekey can leave another node published at the new - // canonical boundary. Move that displaced node to its own current - // canonical key, repeating only if stale boundaries form a short - // collision chain. The fallback scan is confined to such a - // collision; ordinary split/attach publication remains O(log N). - let canonical = self.index.as_deref().expect("topology guard already released"); - let mut route_key = key; - let mut route_node = node; - loop { - let displaced = published.insert(route_key.clone(), route_node.clone()); - published_keys.insert(node_identity(&route_node), route_key.clone()); - - let Some(displaced) = displaced else { - break; - }; - if Arc::ptr_eq(&displaced, &route_node) { - break; - } - published_keys.remove(&node_identity(&displaced)); - - let Some((canonical_key, _)) = canonical - .iter() - .find(|(_, candidate)| Arc::ptr_eq(candidate, &displaced)) - else { - // The displaced route belonged to a node removed from the - // canonical topology by an earlier route-preserving race. - break; - }; - route_key = canonical_key.clone(); - route_node = displaced; + let displaced = self + .published + .as_mut() + .expect("publication snapshot initialized") + .insert(key.clone(), node.clone()); + self.published_keys + .as_mut() + .expect("publication identity map initialized") + .insert(node_identity(&node), key); + if displaced.is_some_and(|old_node| !Arc::ptr_eq(&old_node, &node)) { + // Canonical keys are unique, so a different node cannot + // lawfully occupy this route. Recover once from canonical + // state instead of scanning and chaining while the generation + // remains odd. + self.rebuild_publication(); } } @@ -476,6 +601,7 @@ where T: Borrow, Q: Ord + ?Sized, { + debug_assert!(self.publish, "generic topology removal requires publication"); let removed = self .index .as_deref_mut() @@ -485,21 +611,20 @@ where if self.publish { self.ensure_publication_snapshot(); - let published = self.published.as_mut().expect("publication snapshot initialized"); - let published_keys = self - .published_keys - .as_mut() - .expect("publication identity map initialized"); - let route_key = published_keys - .remove(&node_identity(&removed)) - .expect("published key map must contain a canonical node"); - let route_node = published - .remove::(&route_key) - .expect("published route must contain its recorded boundary"); - assert!( - Arc::ptr_eq(&route_node, &removed), - "published boundary must identify the removed node" - ); + let removed_consistently = { + let published = self.published.as_mut().expect("publication snapshot initialized"); + let published_keys = self + .published_keys + .as_mut() + .expect("publication identity map initialized"); + published_keys + .remove(&node_identity(&removed)) + .and_then(|route_key| published.remove::(&route_key)) + .is_some_and(|route_node| Arc::ptr_eq(&route_node, &removed)) + }; + if !removed_consistently { + self.rebuild_publication(); + } } Some(removed) @@ -530,10 +655,14 @@ where published_keys.len(), "canonical and published identity counts diverged" ); - for node in canonical.values() { + for ((_, canonical_node), (route_key, route_node)) in canonical.iter().zip(published.iter()) { debug_assert!( - published_keys.contains_key(&node_identity(node)), - "canonical node is absent from published identity map" + Arc::ptr_eq(route_node, canonical_node), + "canonical and published node order diverged" + ); + debug_assert!( + published_keys.get(&node_identity(canonical_node)) == Some(route_key), + "published identity key does not match the route index" ); } } @@ -772,8 +901,10 @@ where /// Attaches a persisted topology in one structural publication. /// - /// Nodes must be non-empty, internally sorted, and mutually ordered. The - /// same preconditions as [`Self::attach_node`] apply to every item. + /// Nodes must be non-empty and internally sorted. Their values, together + /// with any nodes already attached to this set, must form mutually ordered + /// non-overlapping ranges. The same preconditions as [`Self::attach_node`] + /// apply to every item. pub fn attach_nodes(&self, nodes: impl IntoIterator) { let mut nodes = nodes.into_iter().peekable(); if nodes.peek().is_none() { @@ -781,6 +912,7 @@ where } let mut index = self.index.write(); + index.repair_last_route_before_attach(); for node in nodes { let node_id = node .max() @@ -3062,6 +3194,83 @@ mod tests { assert_eq!(set.get(&5).map(|value| *value.get()), Some(5)); } + #[test] + fn attach_repairs_a_stale_last_node_boundary() { + let set = BTreeSet::::with_maximum_node_size(8); + set.attach_node(vec![1, 10]); + + // A last-node maximum may remain conservatively published at its old + // high boundary. Incremental restoration above that node must move + // the old route down before installing a new node whose values occupy + // the gap. + assert_eq!(set.remove(&10), Some(10)); + set.attach_node(vec![5, 20]); + + assert_eq!(set.get(&5).map(|value| *value.get()), Some(5)); + assert!(set.contains(&5)); + } + + #[test] + fn attach_repairs_a_stale_low_last_node_boundary() { + let set = BTreeSet::::with_maximum_node_size(8); + set.attach_node(vec![1, 10]); + + // Growing the last node can leave its published boundary below its + // canonical maximum. Once another node is attached, that old route is + // no longer the final fallback and must be repaired as well. + assert!(set.insert(20)); + set.attach_node(vec![25, 30]); + + assert_eq!(set.get(&20).map(|value| *value.get()), Some(20)); + assert_eq!(set.get(&25).map(|value| *value.get()), Some(25)); + } + + #[test] + fn attach_recovers_from_a_missing_published_identity() { + let set = BTreeSet::::with_maximum_node_size(8); + set.attach_node(vec![1, 10]); + let last_identity = { + let index = set.index.read(); + super::node_identity(index.last_key_value().unwrap().1) + }; + assert!(set.index.published_keys.lock().remove(&last_identity).is_some()); + + set.attach_node(vec![20, 30]); + + for value in [1, 10, 20, 30] { + assert_eq!(set.get(&value).map(|found| *found.get()), Some(value)); + } + } + + #[test] + fn attach_boundary_repair_is_a_complete_publication_by_itself() { + let set = BTreeSet::::with_maximum_node_size(8); + set.attach_node(vec![1, 10]); + assert_eq!(set.remove(&10), Some(10)); + + // Model attachment stopping after its preflight repair (for example, + // because user-provided Clone/Ord code panics while reading the first + // incoming node). The repaired identity map and route snapshot must + // still commit together when the guard drops. + { + let mut index = set.index.write(); + index.repair_last_route_before_attach(); + } + set.attach_node(vec![5, 20]); + + assert_eq!(set.get(&5).map(|found| *found.get()), Some(5)); + } + + #[cfg(debug_assertions)] + #[test] + #[should_panic(expected = "generic topology removal requires publication")] + fn publication_must_be_enabled_before_an_opt_out_guard_mutates() { + let set = BTreeSet::::with_maximum_node_size(8); + set.attach_node(vec![1, 10]); + let mut index = set.index.write_rekey(); + index.remove(&10); + } + #[test] fn published_point_routes_match_a_sequential_oracle_under_churn() { let set = BTreeSet::::with_maximum_node_size(4); From 4e04dcc8e0235fe4b3eb8474b4bf89872a17d84c Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 5 Sep 2026 21:25:56 +0700 Subject: [PATCH 4/4] Close low-risk publication audit items --- .github/workflows/miri.yml | 27 ++++++++ CHANGELOG.md | 3 + src/concurrent/set.rs | 128 +++++++++++++++++++++++++++++++++---- 3 files changed, 147 insertions(+), 11 deletions(-) create mode 100644 .github/workflows/miri.yml diff --git a/.github/workflows/miri.yml b/.github/workflows/miri.yml new file mode 100644 index 0000000..9730481 --- /dev/null +++ b/.github/workflows/miri.yml @@ -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_ diff --git a/CHANGELOG.md b/CHANGELOG.md index 464fd11..4d49890 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -35,6 +35,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 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] diff --git a/src/concurrent/set.rs b/src/concurrent/set.rs index c7d1630..40d5065 100644 --- a/src/concurrent/set.rs +++ b/src/concurrent/set.rs @@ -27,6 +27,10 @@ type NodeIndex = BTreeMap>>; // changed boundary; it never copies the full node index. At WorkTable's // default 1,024 rows per node, one 128-route chunk covers roughly 131k rows. const PUBLISHED_ROUTES_PER_CHUNK: usize = 128; +// Leave rebuilt chunks room for subsequent inserts, and merge only below the +// split threshold so alternating insert/remove cannot thrash one boundary. +const PUBLISHED_REBUILD_ROUTES_PER_CHUNK: usize = PUBLISHED_ROUTES_PER_CHUNK * 2 / 3; +const PUBLISHED_ROUTE_MERGE_THRESHOLD: usize = PUBLISHED_ROUTES_PER_CHUNK * 3 / 4; struct PublishedChunk { entries: Vec<(T, Arc>)>, @@ -59,14 +63,14 @@ where T: Ord + Clone, { fn from_canonical(index: &NodeIndex) -> Self { - let mut chunks = Vec::with_capacity(index.len().div_ceil(PUBLISHED_ROUTES_PER_CHUNK)); - let mut entries = Vec::with_capacity(PUBLISHED_ROUTES_PER_CHUNK); + let mut chunks = Vec::with_capacity(index.len().div_ceil(PUBLISHED_REBUILD_ROUTES_PER_CHUNK)); + let mut entries = Vec::with_capacity(PUBLISHED_REBUILD_ROUTES_PER_CHUNK); for (key, node) in index { entries.push((key.clone(), node.clone())); - if entries.len() == PUBLISHED_ROUTES_PER_CHUNK { + if entries.len() == PUBLISHED_REBUILD_ROUTES_PER_CHUNK { chunks.push(Arc::new(PublishedChunk { entries })); - entries = Vec::with_capacity(PUBLISHED_ROUTES_PER_CHUNK); + entries = Vec::with_capacity(PUBLISHED_REBUILD_ROUTES_PER_CHUNK); } } if !entries.is_empty() { @@ -163,12 +167,13 @@ where Q: Ord + ?Sized, { let chunk_index = self.chunk_for(key); - let chunk = self.chunks.get_mut(chunk_index)?; - let chunk = Arc::make_mut(chunk); - let entry_index = chunk + let entry_index = self + .chunks + .get(chunk_index)? .entries .binary_search_by(|(candidate, _)| >::borrow(candidate).cmp(key)) .ok()?; + let chunk = Arc::make_mut(&mut self.chunks[chunk_index]); let (_, removed) = chunk.entries.remove(entry_index); self.len -= 1; @@ -176,7 +181,7 @@ where self.chunks.remove(chunk_index); } else if chunk_index > 0 && self.chunks[chunk_index - 1].entries.len() + self.chunks[chunk_index].entries.len() - <= PUBLISHED_ROUTES_PER_CHUNK + <= PUBLISHED_ROUTE_MERGE_THRESHOLD { let right = self.chunks.remove(chunk_index); Arc::make_mut(&mut self.chunks[chunk_index - 1]) @@ -184,7 +189,7 @@ where .extend(right.entries.iter().cloned()); } else if chunk_index + 1 < self.chunks.len() && self.chunks[chunk_index].entries.len() + self.chunks[chunk_index + 1].entries.len() - <= PUBLISHED_ROUTES_PER_CHUNK + <= PUBLISHED_ROUTE_MERGE_THRESHOLD { let right = self.chunks.remove(chunk_index + 1); Arc::make_mut(&mut self.chunks[chunk_index]) @@ -288,6 +293,14 @@ impl Drop for PublishedIndex { } } +// Publication invariant: every canonical node appears once and in the same +// order in the published route index. At most one route key may differ from +// its canonical key, and only for the canonical last node. That exception is +// safe because point lookup falls back to the published last node above all +// routes, while a stale route below the current maximum still selects that +// same final node. A last-node shrink cannot reorder it before the preceding +// node because node ranges are non-overlapping. Attachment repairs the route +// before it can cease to be the last node. pub(crate) struct Topology { index: RwLock>, // Writer-only reverse lookup from node identity to its current published @@ -1480,8 +1493,9 @@ where let node = first_for_borrowed_bound(&index, Bound::Included(value), self.borrow_order_matches) .or_else(|| index.last_key_value()) .or_else(|| index.first_key_value()) - .map(|(_, node)| node)?; - let node_guard = node.read(); + .map(|(_, node)| node.clone())?; + let node_guard = node.read_arc(); + drop(index); let position = node_guard.try_select(value)?; return node_guard .get_ith(position) @@ -3242,6 +3256,98 @@ mod tests { } } + #[test] + fn insert_recovers_from_a_missing_replaced_node_identity() { + let set = BTreeSet::::with_maximum_node_size(8); + set.attach_node(vec![1, 10]); + let old_node = set.index.read().last_key_value().unwrap().1.clone(); + assert!(set + .index + .published_keys + .lock() + .remove(&super::node_identity(&old_node)) + .is_some()); + + { + let mut index = set.index.write(); + let replaced = index.insert(10, Arc::new(parking_lot::RwLock::new(vec![5, 10]))); + assert!(replaced.is_some_and(|node| Arc::ptr_eq(&node, &old_node))); + } + + assert!(!set.contains(&1)); + assert!(set.contains(&5)); + assert!(set.contains(&10)); + } + + #[test] + fn remove_recovers_from_a_missing_node_identity() { + let set = BTreeSet::::with_maximum_node_size(8); + set.attach_node(vec![1, 10]); + let old_node = set.index.read().last_key_value().unwrap().1.clone(); + assert!(set + .index + .published_keys + .lock() + .remove(&super::node_identity(&old_node)) + .is_some()); + + { + let mut index = set.index.write(); + let removed = index.remove(&10).expect("canonical route exists"); + assert!(Arc::ptr_eq(&removed, &old_node)); + } + + assert!(set.is_empty()); + assert!(!set.contains(&1)); + } + + #[test] + fn missing_published_remove_does_not_clone_a_shared_chunk() { + let mut published = super::PublishedNodeIndex::> { + chunks: Vec::new(), + len: 0, + }; + for key in 0..16 { + published.insert(key, Arc::new(parking_lot::RwLock::new(vec![key]))); + } + let snapshot = published.clone(); + assert!(Arc::ptr_eq(&published.chunks[0], &snapshot.chunks[0])); + + assert!(published.remove(&100).is_none()); + + assert!(Arc::ptr_eq(&published.chunks[0], &snapshot.chunks[0])); + } + + #[test] + fn published_chunks_have_split_merge_hysteresis() { + let mut published = super::PublishedNodeIndex::> { + chunks: Vec::new(), + len: 0, + }; + for key in 0..=128 { + published.insert(key, Arc::new(parking_lot::RwLock::new(vec![key]))); + } + assert_eq!(published.chunks.len(), 2); + + for key in 0..33 { + assert!(published.remove(&key).is_some()); + } + assert_eq!(published.chunks.len(), 1); + + published.insert(0, Arc::new(parking_lot::RwLock::new(vec![0]))); + assert_eq!( + published.chunks.len(), + 1, + "one insert after a merge must not split again" + ); + assert!(published.remove(&0).is_some()); + assert_eq!( + published.chunks.len(), + 1, + "one remove after a merge must not change chunking" + ); + } + #[test] fn attach_boundary_repair_is_a_complete_publication_by_itself() { let set = BTreeSet::::with_maximum_node_size(8);