use granular hierarchical locks in offset_table.rs
This commit is contained in:
@@ -1,13 +1,12 @@
|
||||
use std::cell::UnsafeCell;
|
||||
use std::hash::{Hash, Hasher};
|
||||
use std::ops::{Deref, DerefMut};
|
||||
use std::sync::{Arc, RwLock, RwLockReadGuard};
|
||||
use std::sync::Arc;
|
||||
use std::{fmt, mem, ptr};
|
||||
|
||||
use arcu::atomic::Arcu;
|
||||
use arcu::epoch_counters::GlobalEpochCounterPool;
|
||||
use arcu::rcu_ref::RcuRef;
|
||||
use arcu::Rcu;
|
||||
use parking_lot::RwLock;
|
||||
|
||||
use crate::machine::machine_indices::IndexPtr;
|
||||
use crate::raw_block::RawBlock;
|
||||
@@ -55,7 +54,7 @@ impl<T: RawBlockTraits> From<Arc<ConcurrentOffsetTable<T>>> for OffsetTableImpl<
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: RawBlockTraits> OffsetTableImpl<T> {
|
||||
impl<T: fmt::Debug + RawBlockTraits> OffsetTableImpl<T> {
|
||||
#[inline(always)]
|
||||
pub fn new() -> Self {
|
||||
Self(InnerOffsetTableImpl::Serial(SerialOffsetTable::new()))
|
||||
@@ -70,10 +69,17 @@ impl<T: RawBlockTraits> OffsetTableImpl<T> {
|
||||
};
|
||||
|
||||
let serial_tbl = mem::replace(serial_tbl, empty_serial_tbl);
|
||||
let num_tbl_entries = serial_tbl.block.size() / size_of::<T>();
|
||||
let block = Arcu::new(serial_tbl.block, GlobalEpochCounterPool);
|
||||
|
||||
let growth_lock = RwLock::new(());
|
||||
let concurrent_tbl = Arc::new(ConcurrentOffsetTable { block, growth_lock });
|
||||
let offset_locks: Vec<RwLock<()>> =
|
||||
(0..num_tbl_entries).map(|_| RwLock::new(())).collect();
|
||||
|
||||
let concurrent_tbl = Arc::new(ConcurrentOffsetTable {
|
||||
block,
|
||||
growth_lock: RwLock::new(()),
|
||||
offset_locks: RwLock::new(offset_locks),
|
||||
});
|
||||
|
||||
self.0 = InnerOffsetTableImpl::Concurrent(concurrent_tbl.clone());
|
||||
|
||||
@@ -88,23 +94,48 @@ impl<T: RawBlockTraits> OffsetTableImpl<T> {
|
||||
match &mut self.0 {
|
||||
InnerOffsetTableImpl::Serial(_serial_tbl) => Ok(()),
|
||||
InnerOffsetTableImpl::Concurrent(concurrent_tbl) => {
|
||||
let lock_guard = concurrent_tbl.growth_lock.write().unwrap();
|
||||
let raw_block = concurrent_tbl.block.replace(RawBlock::empty_block());
|
||||
let table_arc = std::mem::replace(
|
||||
concurrent_tbl,
|
||||
Arc::new(ConcurrentOffsetTable {
|
||||
block: Arcu::new(RawBlock::empty_block(), GlobalEpochCounterPool),
|
||||
growth_lock: RwLock::new(()),
|
||||
offset_locks: RwLock::new(vec![]),
|
||||
}),
|
||||
);
|
||||
|
||||
match Arc::try_unwrap(raw_block) {
|
||||
Ok(block) => {
|
||||
drop(lock_guard);
|
||||
self.0 = InnerOffsetTableImpl::Serial(SerialOffsetTable { block });
|
||||
match Arc::try_unwrap(table_arc) {
|
||||
Ok(table) => {
|
||||
// this was the only instance of the concurrent table, as such
|
||||
// at this point no build_with/with_entry{_mut} call can be in-progress/made
|
||||
|
||||
// this shouldn't be able to fail
|
||||
let raw_block =
|
||||
Arc::try_unwrap(table.block.replace(RawBlock::empty_block())).unwrap();
|
||||
self.0 =
|
||||
InnerOffsetTableImpl::Serial(SerialOffsetTable { block: raw_block });
|
||||
Ok(())
|
||||
}
|
||||
Err(_) => Err(()),
|
||||
Err(table_arc) => {
|
||||
// restore the concurrent_tbl
|
||||
*concurrent_tbl = table_arc;
|
||||
Err(())
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[inline]
|
||||
pub fn get_entry(&self, offset: <Self as OffsetTable<T>>::Offset) -> T
|
||||
where
|
||||
Self: OffsetTable<T>,
|
||||
T: Copy,
|
||||
{
|
||||
self.with_entry(offset, |value| *value)
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: RawBlockTraits> Default for OffsetTableImpl<T> {
|
||||
impl<T: fmt::Debug + RawBlockTraits> Default for OffsetTableImpl<T> {
|
||||
fn default() -> Self {
|
||||
Self::new()
|
||||
}
|
||||
@@ -119,6 +150,7 @@ struct SerialOffsetTable<T: RawBlockTraits> {
|
||||
pub struct ConcurrentOffsetTable<T: RawBlockTraits> {
|
||||
block: Arcu<RawBlock<T>, GlobalEpochCounterPool>,
|
||||
growth_lock: RwLock<()>,
|
||||
offset_locks: RwLock<Vec<RwLock<()>>>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
@@ -132,42 +164,24 @@ impl<T: RawBlockTraits> InnerOffsetTableImpl<T> {
|
||||
#[inline(always)]
|
||||
fn build_with(&mut self, value: T) -> usize {
|
||||
match self {
|
||||
Self::Concurrent(concurrent_tbl) => unsafe { concurrent_tbl.build_with(value) },
|
||||
Self::Concurrent(concurrent_tbl) => concurrent_tbl.build_with(value),
|
||||
Self::Serial(serial_tbl) => unsafe { serial_tbl.build_with(value) },
|
||||
}
|
||||
}
|
||||
|
||||
#[inline(always)]
|
||||
fn lookup<'a>(&'a self, offset: usize) -> TablePtr<'a, T> {
|
||||
fn with_entry<R, F: FnOnce(&T) -> R>(&self, offset: usize, f: F) -> R {
|
||||
match self {
|
||||
Self::Concurrent(concurrent_tbl) => TablePtr({
|
||||
let (rcu_ref, guard_lock) = concurrent_tbl.lookup(offset);
|
||||
InnerTablePtr::Concurrent {
|
||||
rcu_ref,
|
||||
guard_lock,
|
||||
}
|
||||
}),
|
||||
Self::Serial(serial_tbl) => unsafe {
|
||||
TablePtr(InnerTablePtr::Serial(serial_tbl.lookup(offset)))
|
||||
},
|
||||
Self::Concurrent(concurrent_tbl) => concurrent_tbl.with_entry(offset, f),
|
||||
Self::Serial(serial_tbl) => f(unsafe { serial_tbl.lookup(offset) }),
|
||||
}
|
||||
}
|
||||
|
||||
#[inline(always)]
|
||||
fn lookup_mut<'a>(&'a mut self, offset: usize) -> TablePtrMut<'a, T> {
|
||||
fn with_entry_mut<R, F: FnOnce(&mut T) -> R>(&mut self, offset: usize, f: F) -> R {
|
||||
match self {
|
||||
InnerOffsetTableImpl::Concurrent(concurrent_tbl) => TablePtrMut({
|
||||
let (rcu_ref, guard_lock) = concurrent_tbl.lookup_mut(offset);
|
||||
InnerTablePtrMut::Concurrent {
|
||||
rcu_ref,
|
||||
guard_lock,
|
||||
}
|
||||
}),
|
||||
InnerOffsetTableImpl::Serial(serial_tbl) => {
|
||||
TablePtrMut(InnerTablePtrMut::Serial(unsafe {
|
||||
serial_tbl.lookup_mut(offset)
|
||||
}))
|
||||
}
|
||||
Self::Concurrent(concurrent_tbl) => concurrent_tbl.with_entry_mut(offset, f),
|
||||
Self::Serial(serial_tbl) => f(unsafe { serial_tbl.lookup_mut(offset) }),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -176,8 +190,9 @@ pub trait OffsetTable<T: RawBlockTraits> {
|
||||
type Offset: Copy + Into<usize>;
|
||||
|
||||
fn build_with(&mut self, value: T) -> Self::Offset;
|
||||
fn lookup<'a>(&'a self, offset: Self::Offset) -> TablePtr<'a, T>;
|
||||
fn lookup_mut<'a>(&'a mut self, offset: Self::Offset) -> TablePtrMut<'a, T>;
|
||||
|
||||
fn with_entry<R, F: FnOnce(&T) -> R>(&self, offset: Self::Offset, f: F) -> R;
|
||||
fn with_entry_mut<R, F: FnOnce(&mut T) -> R>(&mut self, offset: Self::Offset, f: F) -> R;
|
||||
}
|
||||
|
||||
impl OffsetTable<OrderedFloat<f64>> for OffsetTableImpl<OrderedFloat<f64>> {
|
||||
@@ -187,12 +202,18 @@ impl OffsetTable<OrderedFloat<f64>> for OffsetTableImpl<OrderedFloat<f64>> {
|
||||
F64Offset(self.0.build_with(value))
|
||||
}
|
||||
|
||||
fn lookup<'a>(&'a self, offset: F64Offset) -> TablePtr<'a, OrderedFloat<f64>> {
|
||||
self.0.lookup(offset.into())
|
||||
#[inline]
|
||||
fn with_entry<R, F: FnOnce(&OrderedFloat<f64>) -> R>(&self, offset: F64Offset, f: F) -> R {
|
||||
self.0.with_entry(offset.into(), f)
|
||||
}
|
||||
|
||||
fn lookup_mut<'a>(&'a mut self, offset: F64Offset) -> TablePtrMut<'a, OrderedFloat<f64>> {
|
||||
self.0.lookup_mut(offset.into())
|
||||
#[inline]
|
||||
fn with_entry_mut<R, F: FnOnce(&mut OrderedFloat<f64>) -> R>(
|
||||
&mut self,
|
||||
offset: F64Offset,
|
||||
f: F,
|
||||
) -> R {
|
||||
self.0.with_entry_mut(offset.into(), f)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -203,12 +224,18 @@ impl OffsetTable<IndexPtr> for OffsetTableImpl<IndexPtr> {
|
||||
CodeIndexOffset(self.0.build_with(value))
|
||||
}
|
||||
|
||||
fn lookup<'a>(&'a self, offset: CodeIndexOffset) -> TablePtr<'a, IndexPtr> {
|
||||
self.0.lookup(offset.into())
|
||||
#[inline]
|
||||
fn with_entry<R, F: FnOnce(&IndexPtr) -> R>(&self, offset: CodeIndexOffset, f: F) -> R {
|
||||
self.0.with_entry(offset.into(), f)
|
||||
}
|
||||
|
||||
fn lookup_mut<'a>(&'a mut self, offset: CodeIndexOffset) -> TablePtrMut<'a, IndexPtr> {
|
||||
self.0.lookup_mut(offset.into())
|
||||
#[inline]
|
||||
fn with_entry_mut<R, F: FnOnce(&mut IndexPtr) -> R>(
|
||||
&mut self,
|
||||
offset: CodeIndexOffset,
|
||||
f: F,
|
||||
) -> R {
|
||||
self.0.with_entry_mut(offset.into(), f)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -251,8 +278,8 @@ impl<T: RawBlockTraits> SerialOffsetTable<T> {
|
||||
|
||||
impl<T: RawBlockTraits> ConcurrentOffsetTable<T> {
|
||||
#[allow(clippy::missing_safety_doc)]
|
||||
unsafe fn build_with(&self, value: T) -> usize {
|
||||
let update_guard = self.growth_lock.write().unwrap();
|
||||
fn build_with(&self, value: T) -> usize {
|
||||
let growth_lock = self.growth_lock.write();
|
||||
|
||||
// we don't have an index table for lookups as AtomTable does so
|
||||
// just get the epoch after we take the upgrade lock
|
||||
@@ -260,10 +287,10 @@ impl<T: RawBlockTraits> ConcurrentOffsetTable<T> {
|
||||
let mut ptr;
|
||||
|
||||
loop {
|
||||
ptr = block_epoch.alloc(mem::size_of::<T>());
|
||||
ptr = unsafe { block_epoch.alloc(mem::size_of::<T>()) };
|
||||
|
||||
if ptr.is_null() {
|
||||
let new_block = block_epoch.grow_new().unwrap();
|
||||
let new_block = unsafe { block_epoch.grow_new().unwrap() };
|
||||
self.block.replace(new_block);
|
||||
block_epoch = self.block.read();
|
||||
} else {
|
||||
@@ -271,35 +298,46 @@ impl<T: RawBlockTraits> ConcurrentOffsetTable<T> {
|
||||
}
|
||||
}
|
||||
|
||||
ptr::write(ptr as *mut T, value);
|
||||
let new_tbl_sz = block_epoch.size() / size_of::<T>();
|
||||
let mut offset_locks = self.offset_locks.write();
|
||||
|
||||
offset_locks.resize_with(new_tbl_sz, || RwLock::new(()));
|
||||
|
||||
unsafe {
|
||||
ptr::write(ptr as *mut T, value);
|
||||
}
|
||||
|
||||
let value = ptr.addr() - block_epoch.base.addr();
|
||||
|
||||
// AtomTable would have to update the index table at this point
|
||||
// explicit drop to ensure we don't accidentally drop it early
|
||||
drop(update_guard);
|
||||
drop(offset_locks);
|
||||
drop(growth_lock);
|
||||
|
||||
value
|
||||
}
|
||||
|
||||
#[inline]
|
||||
fn lookup<'a>(&'a self, offset: usize) -> (RcuRef<RawBlock<T>, T>, RwLockReadGuard<'a, ()>) {
|
||||
let growth_lock_guard = self.growth_lock.read().unwrap();
|
||||
fn with_entry<R, F: FnOnce(&T) -> R>(&self, offset: usize, f: F) -> R {
|
||||
let outer_offset_lock = self.offset_locks.read();
|
||||
let inner_offset_lock = outer_offset_lock[offset / size_of::<T>()].read();
|
||||
|
||||
let rcu_ref = RcuRef::try_map(self.block.read(), |raw_block| unsafe {
|
||||
raw_block.base.add(offset).cast::<T>().as_ref()
|
||||
})
|
||||
.expect("The offset should result in a non-null pointer");
|
||||
.expect("offset valid");
|
||||
|
||||
(rcu_ref, growth_lock_guard)
|
||||
let result = f(&*rcu_ref);
|
||||
|
||||
drop(inner_offset_lock);
|
||||
drop(outer_offset_lock);
|
||||
|
||||
result
|
||||
}
|
||||
|
||||
#[inline]
|
||||
fn lookup_mut<'a>(
|
||||
&'a self,
|
||||
offset: usize,
|
||||
) -> (RcuRef<RawBlock<T>, UnsafeCell<T>>, RwLockReadGuard<'a, ()>) {
|
||||
let growth_lock_guard = self.growth_lock.read().unwrap();
|
||||
fn with_entry_mut<R, F: FnOnce(&mut T) -> R>(&self, offset: usize, f: F) -> R {
|
||||
let growth_lock = self.growth_lock.read();
|
||||
let outer_offset_lock = self.offset_locks.read();
|
||||
let inner_offset_lock = outer_offset_lock[offset / size_of::<T>()].write();
|
||||
|
||||
let rcu_ref = RcuRef::try_map(self.block.read(), |raw_block| unsafe {
|
||||
raw_block
|
||||
@@ -309,9 +347,15 @@ impl<T: RawBlockTraits> ConcurrentOffsetTable<T> {
|
||||
.cast::<UnsafeCell<T>>()
|
||||
.as_ref()
|
||||
})
|
||||
.expect("The offset should result in a non-null pointer");
|
||||
.expect("offset valid");
|
||||
|
||||
(rcu_ref, growth_lock_guard)
|
||||
let result = f(unsafe { &mut *rcu_ref.get().as_mut().unwrap() });
|
||||
|
||||
drop(inner_offset_lock);
|
||||
drop(outer_offset_lock);
|
||||
drop(growth_lock);
|
||||
|
||||
result
|
||||
}
|
||||
}
|
||||
|
||||
@@ -358,64 +402,6 @@ impl CodeIndexOffset {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct TablePtr<'a, T: RawBlockTraits>(InnerTablePtr<'a, T>);
|
||||
|
||||
#[derive(Debug)]
|
||||
enum InnerTablePtr<'a, T: RawBlockTraits> {
|
||||
Concurrent {
|
||||
rcu_ref: RcuRef<RawBlock<T>, T>,
|
||||
#[allow(dead_code)]
|
||||
guard_lock: RwLockReadGuard<'a, ()>,
|
||||
},
|
||||
Serial(&'a T),
|
||||
}
|
||||
|
||||
impl<T: PartialEq + RawBlockTraits> PartialEq for TablePtr<'_, T> {
|
||||
fn eq(&self, other: &TablePtr<'_, T>) -> bool {
|
||||
self.deref() == other.deref()
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: Eq + RawBlockTraits> Eq for TablePtr<'_, T> {}
|
||||
|
||||
impl<T: PartialOrd + Ord + RawBlockTraits> PartialOrd for TablePtr<'_, T> {
|
||||
fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
|
||||
Some(self.cmp(other))
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: Ord + RawBlockTraits> Ord for TablePtr<'_, T> {
|
||||
fn cmp(&self, other: &Self) -> std::cmp::Ordering {
|
||||
(**self).cmp(&**other)
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: Hash + RawBlockTraits> Hash for TablePtr<'_, T> {
|
||||
#[inline(always)]
|
||||
fn hash<H: Hasher>(&self, hasher: &mut H) {
|
||||
(self as &T).hash(hasher)
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: fmt::Display + RawBlockTraits> fmt::Display for TablePtr<'_, T> {
|
||||
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
|
||||
write!(f, "{}", self as &T)
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: RawBlockTraits> Deref for TablePtr<'_, T> {
|
||||
type Target = T;
|
||||
|
||||
#[inline]
|
||||
fn deref(&self) -> &Self::Target {
|
||||
match &self.0 {
|
||||
InnerTablePtr::Concurrent { rcu_ref, .. } => rcu_ref,
|
||||
InnerTablePtr::Serial(ref_mut) => ref_mut,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Display for CodeIndexOffset {
|
||||
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
|
||||
write!(f, "CodeIndexOffset({})", self.0)
|
||||
@@ -434,97 +420,3 @@ impl fmt::Display for F64Offset {
|
||||
write!(f, "F64Offset({})", self.0)
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct TablePtrMut<'a, T: RawBlockTraits>(InnerTablePtrMut<'a, T>);
|
||||
|
||||
#[derive(Debug)]
|
||||
enum InnerTablePtrMut<'a, T: RawBlockTraits> {
|
||||
Concurrent {
|
||||
rcu_ref: RcuRef<RawBlock<T>, UnsafeCell<T>>,
|
||||
#[allow(dead_code)]
|
||||
guard_lock: RwLockReadGuard<'a, ()>,
|
||||
},
|
||||
Serial(&'a mut T),
|
||||
}
|
||||
|
||||
impl<T: PartialEq + RawBlockTraits> PartialEq for TablePtrMut<'_, T> {
|
||||
fn eq(&self, other: &TablePtrMut<'_, T>) -> bool {
|
||||
self.deref() == other.deref()
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: Eq + RawBlockTraits> Eq for TablePtrMut<'_, T> {}
|
||||
|
||||
impl<T: PartialOrd + Ord + RawBlockTraits> PartialOrd for TablePtrMut<'_, T> {
|
||||
fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
|
||||
Some(self.cmp(other))
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: Ord + RawBlockTraits> Ord for TablePtrMut<'_, T> {
|
||||
fn cmp(&self, other: &Self) -> std::cmp::Ordering {
|
||||
(**self).cmp(&**other)
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: Hash + RawBlockTraits> Hash for TablePtrMut<'_, T> {
|
||||
#[inline(always)]
|
||||
fn hash<H: Hasher>(&self, hasher: &mut H) {
|
||||
(self as &T).hash(hasher)
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: fmt::Display + RawBlockTraits> fmt::Display for TablePtrMut<'_, T> {
|
||||
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
|
||||
write!(f, "{}", self as &T)
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: RawBlockTraits> Deref for TablePtrMut<'_, T> {
|
||||
type Target = T;
|
||||
|
||||
#[inline]
|
||||
fn deref(&self) -> &Self::Target {
|
||||
match &self.0 {
|
||||
InnerTablePtrMut::Concurrent { rcu_ref, .. } => unsafe {
|
||||
rcu_ref.get().as_ref().unwrap()
|
||||
},
|
||||
InnerTablePtrMut::Serial(ref_mut) => ref_mut,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: RawBlockTraits> DerefMut for TablePtrMut<'_, T> {
|
||||
#[inline]
|
||||
fn deref_mut(&mut self) -> &mut Self::Target {
|
||||
match &mut self.0 {
|
||||
InnerTablePtrMut::Concurrent { rcu_ref, .. } => unsafe {
|
||||
&mut *rcu_ref.get().as_mut().unwrap()
|
||||
},
|
||||
InnerTablePtrMut::Serial(ref_mut) => ref_mut,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl TablePtrMut<'_, IndexPtr> {
|
||||
#[inline]
|
||||
pub fn set(&mut self, val: IndexPtr) {
|
||||
match &mut self.0 {
|
||||
InnerTablePtrMut::Concurrent { rcu_ref, .. } => unsafe {
|
||||
*rcu_ref.get() = val;
|
||||
},
|
||||
InnerTablePtrMut::Serial(ref_mut) => {
|
||||
**ref_mut = val;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[inline]
|
||||
pub fn replace(&mut self, val: IndexPtr) -> IndexPtr {
|
||||
match &mut self.0 {
|
||||
InnerTablePtrMut::Concurrent { rcu_ref, .. } => unsafe { rcu_ref.get().replace(val) },
|
||||
InnerTablePtrMut::Serial(ref_mut) => mem::replace(*ref_mut, val),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user