use parking_lot::RwLockWriteGuard;
use wdev::Device;
use wkv::RangeIndexManager as Engine;
use crate::storage::session::storage_session::StorageSession;
pub type ExclusiveRangeIndexLock<'a> = RwLockWriteGuard<'a, ()>;
pub struct RangeIndexManager_Locking;
impl RangeIndexManager_Locking {
pub fn acquire_exclusive_for_delete(
engine: &Engine,
key_hash: u64,
) -> ExclusiveRangeIndexLock<'_> {
engine.locks().write(key_hash)
}
pub fn promote_to_tail<'a, D: Device>(session: &StorageSession<'a, D>) {
let _promoted = session.promote_range_index_to_tail();
}
}
impl RangeIndexManager_Locking {
pub fn acquire_exclusive_for_key<'a>(
engine: &'a Engine,
key: &[u8],
) -> ExclusiveRangeIndexLock<'a> {
Self::acquire_exclusive_for_delete(engine, Engine::key_hash_of(key))
}
}
#[cfg(test)]
mod tests {
use std::{
sync::{Arc, mpsc},
thread,
time::Duration,
};
use tempfile::tempdir;
use super::*;
fn engine() -> (tempfile::TempDir, Arc<Engine>) {
let dir = tempdir().unwrap();
let e = Arc::new(Engine::new(dir.path().join("ri"), dir.path().join("cpr")));
(dir, e)
}
#[test]
fn exclusive_lock_serializes_cross_thread() {
let (_dir, engine) = engine();
let key_hash = Engine::key_hash_of(b"del-key");
{
let _guard = RangeIndexManager_Locking::acquire_exclusive_for_delete(&engine, key_hash);
let engine2 = Arc::clone(&engine);
let (tx, rx) = mpsc::channel();
let h = thread::spawn(move || {
let g = RangeIndexManager_Locking::acquire_exclusive_for_delete(&engine2, key_hash);
tx.send(()).unwrap();
drop(g);
});
assert!(rx.recv_timeout(Duration::from_millis(100)).is_err());
drop(_guard);
rx.recv_timeout(Duration::from_secs(5)).unwrap();
h.join().unwrap();
}
}
#[test]
fn acquire_exclusive_for_key_derives_same_stripe() {
let (_dir, engine) = engine();
let g1 =
RangeIndexManager_Locking::acquire_exclusive_for_delete(&engine, Engine::key_hash_of(b"k"));
drop(g1);
let g2 = RangeIndexManager_Locking::acquire_exclusive_for_key(&engine, b"k");
drop(g2);
}
}