use super::lock_entry::LockMode;
use super::lock_key::TxnId;
use super::manager::LockManager;
impl LockManager {
pub fn reap_expired_shared(&mut self, epoch_threshold: u64) -> Vec<TxnId> {
let expired: Vec<TxnId> = self
.held_locks
.iter()
.filter(|(owner, keys)| {
owner.is_reservation()
&& owner.epoch < epoch_threshold
&& keys.iter().all(|k| {
self.table
.get(k)
.is_some_and(|e| e.mode == LockMode::Shared)
})
})
.map(|(owner, _)| *owner)
.collect();
let mut promoted = Vec::new();
for owner in expired {
tracing::debug!(
epoch = owner.epoch,
position = owner.position,
"calvin: reaping lease-expired shared reservation"
);
promoted.extend(self.release(owner));
}
promoted
}
}
#[cfg(test)]
mod tests {
use std::collections::BTreeSet;
use std::sync::Arc;
use super::*;
use crate::control::cluster::calvin::scheduler::lock::lock_entry::AcquireOutcome;
use crate::control::cluster::calvin::scheduler::lock::lock_key::LockKey;
fn key(name: &str) -> LockKey {
LockKey::Surrogate {
collection: Arc::from(name),
surrogate: 1,
}
}
fn keyset(names: &[&str]) -> BTreeSet<LockKey> {
names.iter().map(|n| key(n)).collect()
}
fn txn(epoch: u64, pos: u32) -> TxnId {
TxnId::new(epoch, pos)
}
fn resv(epoch: u64, pos_offset: u32) -> TxnId {
TxnId::new(epoch, TxnId::RESERVATION_POSITION_BAND + pos_offset)
}
#[test]
fn reap_releases_expired_shared_reservation() {
let mut lm = LockManager::new();
let r = resv(1, 0);
assert_eq!(lm.acquire_shared(r, key("k")), AcquireOutcome::Ready);
assert_eq!(lm.lock_count(), 1);
let promoted = lm.reap_expired_shared(100);
assert!(promoted.is_empty(), "no waiter queued behind the key");
assert_eq!(lm.lock_count(), 0, "the reservation's key is freed");
assert_eq!(lm.holder_count(), 0, "the reservation is released");
}
#[test]
fn reap_ignores_within_lease() {
let mut lm = LockManager::new();
let r = resv(90, 0);
assert_eq!(lm.acquire_shared(r, key("k")), AcquireOutcome::Ready);
let promoted = lm.reap_expired_shared(50);
assert!(promoted.is_empty());
assert_eq!(lm.lock_count(), 1, "still within the lease, not reaped");
assert!(lm.table.get(&key("k")).unwrap().holders.contains(&r));
}
#[test]
fn reap_ignores_non_reservation_owner() {
let mut lm = LockManager::new();
let t = txn(1, 0);
assert_eq!(lm.acquire(t, keyset(&["k"])), AcquireOutcome::Ready);
let promoted = lm.reap_expired_shared(100);
assert!(promoted.is_empty());
assert_eq!(
lm.lock_count(),
1,
"a real txn's exclusive lock is never reaped"
);
assert!(lm.table.get(&key("k")).unwrap().holders.contains(&t));
}
#[test]
fn reap_skips_self_upgraded_reservation() {
let mut lm = LockManager::new();
let r = resv(1, 0);
assert_eq!(lm.acquire_shared(r, key("k")), AcquireOutcome::Ready);
assert_eq!(lm.acquire(r, keyset(&["k"])), AcquireOutcome::Ready);
let promoted = lm.reap_expired_shared(100);
assert!(promoted.is_empty());
let entry = lm.table.get(&key("k")).unwrap();
assert_eq!(
entry.mode,
LockMode::Exclusive,
"mid-commit reservation is left alone"
);
assert_eq!(entry.holders.len(), 1);
assert_eq!(entry.holders[0], r);
}
#[test]
fn reap_promotes_waiter() {
let mut lm = LockManager::new();
let r = resv(1, 0);
let writer = txn(5, 0);
assert_eq!(lm.acquire_shared(r, key("k")), AcquireOutcome::Ready);
assert_eq!(lm.acquire(writer, keyset(&["k"])), AcquireOutcome::Blocked);
assert!(
lm.table
.get(&key("k"))
.unwrap()
.waiters
.iter()
.any(|(w, _)| *w == writer)
);
let promoted = lm.reap_expired_shared(100);
assert_eq!(
promoted,
vec![writer],
"reaping the stuck reservation unblocks the younger writer"
);
let entry = lm.table.get(&key("k")).unwrap();
assert_eq!(entry.mode, LockMode::Exclusive);
assert_eq!(entry.holders.len(), 1);
assert_eq!(entry.holders[0], writer);
}
}