use std::collections::HashMap;
use std::sync::atomic::{AtomicU8, AtomicU64, Ordering};
use std::time::Duration;
use fsqlite_types::{PageNumber, PageNumberBuildHasher, TxnId};
use parking_lot::Mutex;
use crate::cache_aligned::CacheAligned;
pub const MAX_FC_SLOTS: usize = 64;
const LARGE_BATCH_LOG_THRESHOLD: u32 = 8;
const FC_HANDOFF_BASE_SPINS: u32 = 64;
const FC_HANDOFF_MAX_SPINS: u32 = 2_048;
const FC_HANDOFF_PARK_EVERY: u32 = 4;
const FC_HANDOFF_MAX_PARK: Duration = Duration::from_micros(50);
static FC_SLOT_FULL_EVENTS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
static FC_SLOT_FULL_PARK_NS_TOTAL: std::sync::atomic::AtomicU64 =
std::sync::atomic::AtomicU64::new(0);
fn record_fc_slot_full_event() {
FC_SLOT_FULL_EVENTS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
#[cfg(not(target_arch = "wasm32"))]
fn record_fc_slot_full_park_ns(ns: u64) {
FC_SLOT_FULL_PARK_NS_TOTAL.fetch_add(ns, std::sync::atomic::Ordering::Relaxed);
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
pub struct FcSlotFullMetrics {
pub slot_full_events: u64,
pub park_ns_total: u64,
}
#[must_use]
pub fn fc_slot_full_metrics() -> FcSlotFullMetrics {
FcSlotFullMetrics {
slot_full_events: FC_SLOT_FULL_EVENTS.load(std::sync::atomic::Ordering::Relaxed),
park_ns_total: FC_SLOT_FULL_PARK_NS_TOTAL.load(std::sync::atomic::Ordering::Relaxed),
}
}
pub fn reset_fc_slot_full_metrics() {
FC_SLOT_FULL_EVENTS.store(0, std::sync::atomic::Ordering::Relaxed);
FC_SLOT_FULL_PARK_NS_TOTAL.store(0, std::sync::atomic::Ordering::Relaxed);
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct FcHandoffWait {
attempt: u32,
spin_loops: u32,
park_timeout: Duration,
}
const fn fc_handoff_spin_loops(attempt: u32) -> u32 {
let growth = attempt.saturating_sub(1);
let shift = if growth > 5 { 5 } else { growth };
let spins = FC_HANDOFF_BASE_SPINS << shift;
if spins > FC_HANDOFF_MAX_SPINS {
FC_HANDOFF_MAX_SPINS
} else {
spins
}
}
const fn fc_handoff_should_park(attempt: u32) -> bool {
attempt >= FC_HANDOFF_PARK_EVERY && attempt.is_multiple_of(FC_HANDOFF_PARK_EVERY)
}
const fn fc_handoff_wait(attempt: u32) -> FcHandoffWait {
FcHandoffWait {
attempt,
spin_loops: fc_handoff_spin_loops(attempt),
park_timeout: if fc_handoff_should_park(attempt) {
FC_HANDOFF_MAX_PARK
} else {
Duration::ZERO
},
}
}
fn perform_fc_handoff_spin(wait: FcHandoffWait) {
for _ in 0..wait.spin_loops {
std::hint::spin_loop();
}
}
const OWNER_VACANT: u64 = 0;
const SLOT_IDLE: u8 = 0;
const SLOT_PUBLISHED: u8 = 1;
const SLOT_READY: u8 = 2;
const SLOT_CANCELLED: u8 = 3;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum FcOp {
TryAcquire { page: PageNumber, txn: TxnId },
Release { page: PageNumber, txn: TxnId },
Holder { page: PageNumber },
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum FcOutcome {
Acquired,
Held(TxnId),
Released,
NotHeld,
Unlocked,
HolderIs(TxnId),
}
struct FcSlotInner {
request: Option<FcOp>,
outcome: Option<FcOutcome>,
}
impl FcSlotInner {
const fn new() -> Self {
Self {
request: None,
outcome: None,
}
}
}
struct FcSlot {
owner: AtomicU64,
state: AtomicU8,
inner: Mutex<FcSlotInner>,
parked_thread: Mutex<Option<std::thread::Thread>>,
}
impl FcSlot {
fn new() -> Self {
Self {
owner: AtomicU64::new(OWNER_VACANT),
state: AtomicU8::new(SLOT_IDLE),
inner: Mutex::new(FcSlotInner::new()),
parked_thread: Mutex::new(None),
}
}
}
fn register_slot_waiter(slot: &FcSlot) {
*slot.parked_thread.lock() = Some(std::thread::current());
}
fn clear_slot_waiter(slot: &FcSlot) {
let _ = slot.parked_thread.lock().take();
}
fn unpark_slot_waiter(slot: &FcSlot) {
let waiter = slot.parked_thread.lock().take();
if let Some(waiter) = waiter {
waiter.unpark();
}
}
fn park_current_thread_for_slot(slot: &FcSlot, timeout: Duration) {
if timeout.is_zero() {
return;
}
register_slot_waiter(slot);
if slot.state.load(Ordering::Acquire) == SLOT_READY {
clear_slot_waiter(slot);
return;
}
#[cfg(not(target_arch = "wasm32"))]
std::thread::park_timeout(timeout);
#[cfg(target_arch = "wasm32")]
{
let _ = timeout;
clear_slot_waiter(slot);
}
}
struct FcBackingMap {
inner: Mutex<HashMap<PageNumber, TxnId, PageNumberBuildHasher>>,
}
impl FcBackingMap {
fn new() -> Self {
Self {
inner: Mutex::new(HashMap::with_hasher(PageNumberBuildHasher::default())),
}
}
}
pub struct FcPageLockShard {
combiner_lock: Mutex<()>,
slots: Box<[CacheAligned<FcSlot>; MAX_FC_SLOTS]>,
map: FcBackingMap,
scan_counter: AtomicU64,
shard_index: u32,
}
impl FcPageLockShard {
#[must_use]
pub fn new(shard_index: u32) -> Self {
Self {
combiner_lock: Mutex::new(()),
slots: Box::new(std::array::from_fn(|_| CacheAligned::new(FcSlot::new()))),
map: FcBackingMap::new(),
scan_counter: AtomicU64::new(0),
shard_index,
}
}
pub fn try_acquire(&self, page: PageNumber, txn: TxnId) -> Result<(), TxnId> {
match self.submit(FcOp::TryAcquire { page, txn }) {
FcOutcome::Acquired => Ok(()),
FcOutcome::Held(h) => Err(h),
other => {
assert!(
matches!(&other, FcOutcome::Acquired | FcOutcome::Held(_)),
"try_acquire expected Acquired|Held, got {other:?}"
);
Err(txn)
}
}
}
pub fn release(&self, page: PageNumber, txn: TxnId) -> bool {
match self.submit(FcOp::Release { page, txn }) {
FcOutcome::Released => true,
FcOutcome::NotHeld => false,
other => {
assert!(
matches!(&other, FcOutcome::Released | FcOutcome::NotHeld),
"release expected Released|NotHeld, got {other:?}"
);
false
}
}
}
#[must_use]
pub fn holder(&self, page: PageNumber) -> Option<TxnId> {
match self.submit(FcOp::Holder { page }) {
FcOutcome::Unlocked => None,
FcOutcome::HolderIs(h) => Some(h),
other => {
assert!(
matches!(&other, FcOutcome::Unlocked | FcOutcome::HolderIs(_)),
"holder expected Unlocked|HolderIs, got {other:?}"
);
None
}
}
}
pub fn release_all(&self, txn: TxnId) -> usize {
let _guard = self.combiner_lock.lock();
self.drain_locked();
let mut map = self.map.inner.lock();
let before = map.len();
map.retain(|_, &mut v| v != txn);
before - map.len()
}
#[must_use]
pub fn len(&self) -> usize {
let _guard = self.combiner_lock.lock();
self.map.inner.lock().len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn for_each<F: FnMut(PageNumber, TxnId)>(&self, mut f: F) {
let _guard = self.combiner_lock.lock();
self.drain_locked();
let map = self.map.inner.lock();
for (&p, &t) in map.iter() {
f(p, t);
}
}
pub fn retain<P>(&self, mut predicate: P) -> usize
where
P: FnMut(PageNumber, TxnId) -> bool,
{
let _guard = self.combiner_lock.lock();
self.drain_locked();
let mut map = self.map.inner.lock();
let before = map.len();
map.retain(|&page, &mut txn| predicate(page, txn));
before - map.len()
}
fn submit(&self, op: FcOp) -> FcOutcome {
if let Some(guard) = self.combiner_lock.try_lock() {
self.drain_locked();
let outcome = self.execute_locked(op);
drop(guard);
return outcome;
}
self.submit_via_slot(op)
}
fn submit_via_slot(&self, op: FcOp) -> FcOutcome {
let slot_idx = self.acquire_slot();
let slot = &self.slots[slot_idx];
{
let mut inner = slot.inner.lock();
inner.request = Some(op);
inner.outcome = None;
}
slot.state.store(SLOT_PUBLISHED, Ordering::Release);
if let Some(guard) = self.combiner_lock.try_lock() {
self.drain_locked();
drop(guard);
}
let mut wait_attempt: u32 = 0;
loop {
let st = slot.state.load(Ordering::Acquire);
if st == SLOT_READY {
let outcome = {
let mut inner = slot.inner.lock();
inner.request = None;
inner.outcome.take()
};
clear_slot_waiter(slot);
slot.state.store(SLOT_IDLE, Ordering::Release);
slot.owner.store(OWNER_VACANT, Ordering::Release);
return outcome.expect("combiner set SLOT_READY without outcome");
}
wait_attempt = wait_attempt.saturating_add(1);
let wait = fc_handoff_wait(wait_attempt);
perform_fc_handoff_spin(wait);
if let Some(guard) = self.combiner_lock.try_lock() {
self.drain_locked();
drop(guard);
} else {
park_current_thread_for_slot(slot, wait.park_timeout);
}
}
}
fn acquire_slot(&self) -> usize {
let tid = thread_id_hash();
let start = (tid as usize) % MAX_FC_SLOTS;
let mut wait_attempt: u32 = 0;
loop {
for offset in 0..MAX_FC_SLOTS {
let idx = (start + offset) % MAX_FC_SLOTS;
let slot = &self.slots[idx];
if slot
.owner
.compare_exchange(OWNER_VACANT, tid, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
{
return idx;
}
}
record_fc_slot_full_event();
wait_attempt = wait_attempt.saturating_add(1);
let wait = fc_handoff_wait(wait_attempt);
perform_fc_handoff_spin(wait);
#[cfg(not(target_arch = "wasm32"))]
if !wait.park_timeout.is_zero() {
let park_started = std::time::Instant::now();
std::thread::park_timeout(wait.park_timeout);
record_fc_slot_full_park_ns(
u64::try_from(park_started.elapsed().as_nanos()).unwrap_or(u64::MAX),
);
}
}
}
fn drain_locked(&self) {
let scan = self.scan_counter.fetch_add(1, Ordering::Relaxed);
let mut batch: u32 = 0;
for slot in self.slots.iter() {
let st = slot.state.load(Ordering::Acquire);
match st {
SLOT_PUBLISHED => {
let op = {
let inner = slot.inner.lock();
inner.request
};
let Some(op) = op else {
slot.state.store(SLOT_IDLE, Ordering::Release);
continue;
};
let outcome = self.execute_locked(op);
{
let mut inner = slot.inner.lock();
inner.outcome = Some(outcome);
}
slot.state.store(SLOT_READY, Ordering::Release);
unpark_slot_waiter(slot);
batch += 1;
}
SLOT_CANCELLED => {
{
let mut inner = slot.inner.lock();
inner.request = None;
inner.outcome = None;
}
unpark_slot_waiter(slot);
slot.state.store(SLOT_IDLE, Ordering::Release);
slot.owner.store(OWNER_VACANT, Ordering::Release);
}
_ => {
}
}
}
if batch >= LARGE_BATCH_LOG_THRESHOLD {
tracing::info!(
target: "fsqlite.mvcc.page_lock_fc",
shard = self.shard_index,
scan,
batch_size = batch,
"flat_combine_large_batch"
);
} else if batch > 0 {
tracing::debug!(
target: "fsqlite.mvcc.page_lock_fc",
shard = self.shard_index,
scan,
batch_size = batch,
"flat_combine_batch"
);
}
}
fn execute_locked(&self, op: FcOp) -> FcOutcome {
let mut map = self.map.inner.lock();
match op {
FcOp::TryAcquire { page, txn } => {
if let Some(&holder) = map.get(&page) {
if holder == txn {
FcOutcome::Acquired
} else {
FcOutcome::Held(holder)
}
} else {
map.insert(page, txn);
FcOutcome::Acquired
}
}
FcOp::Release { page, txn } => {
if map.get(&page) == Some(&txn) {
map.remove(&page);
FcOutcome::Released
} else {
FcOutcome::NotHeld
}
}
FcOp::Holder { page } => match map.get(&page).copied() {
Some(h) => FcOutcome::HolderIs(h),
None => FcOutcome::Unlocked,
},
}
}
}
impl std::fmt::Debug for FcPageLockShard {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("FcPageLockShard")
.field("shard_index", &self.shard_index)
.field("scan_counter", &self.scan_counter.load(Ordering::Relaxed))
.finish_non_exhaustive()
}
}
fn thread_id_hash() -> u64 {
let id = std::thread::current().id();
let s = format!("{id:?}");
let mut h: u64 = 0x9E37_79B9_7F4A_7C15;
for b in s.as_bytes() {
h = h.wrapping_mul(0x100_0000_01B3).wrapping_add(u64::from(*b));
}
if h == 0 { 1 } else { h }
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::{Arc, Barrier};
use std::thread;
fn page(n: u32) -> PageNumber {
PageNumber::new(n).expect("non-zero")
}
fn txn(n: u64) -> TxnId {
TxnId::new(n).expect("non-zero txn id")
}
#[test]
fn handoff_spin_loops_grow_then_cap() {
assert_eq!(fc_handoff_spin_loops(1), FC_HANDOFF_BASE_SPINS);
assert_eq!(fc_handoff_spin_loops(2), FC_HANDOFF_BASE_SPINS * 2);
assert_eq!(fc_handoff_spin_loops(3), FC_HANDOFF_BASE_SPINS * 4);
assert_eq!(fc_handoff_spin_loops(6), FC_HANDOFF_MAX_SPINS);
assert_eq!(fc_handoff_spin_loops(32), FC_HANDOFF_MAX_SPINS);
}
#[test]
fn handoff_parks_only_on_bounded_cadence() {
for attempt in 1..FC_HANDOFF_PARK_EVERY {
assert!(
!fc_handoff_should_park(attempt),
"attempt {attempt} should stay on CPU"
);
}
assert!(fc_handoff_should_park(FC_HANDOFF_PARK_EVERY));
assert!(!fc_handoff_should_park(FC_HANDOFF_PARK_EVERY + 1));
assert!(fc_handoff_should_park(FC_HANDOFF_PARK_EVERY * 2));
let wait = fc_handoff_wait(FC_HANDOFF_PARK_EVERY);
assert_eq!(wait.spin_loops, FC_HANDOFF_BASE_SPINS << 3);
assert_eq!(wait.park_timeout, FC_HANDOFF_MAX_PARK);
}
#[test]
fn published_slot_unparks_registered_waiter() {
let slot = Arc::new(FcSlot::new());
slot.state.store(SLOT_PUBLISHED, Ordering::Release);
let (registered_tx, registered_rx) = std::sync::mpsc::channel();
let waiter_slot = Arc::clone(&slot);
let waiter = thread::spawn(move || {
register_slot_waiter(&waiter_slot);
registered_tx.send(()).unwrap();
std::thread::park_timeout(Duration::from_secs(1));
waiter_slot.state.load(Ordering::Acquire)
});
registered_rx
.recv_timeout(Duration::from_secs(1))
.expect("waiter should register before publish");
slot.state.store(SLOT_READY, Ordering::Release);
unpark_slot_waiter(&slot);
assert_eq!(waiter.join().unwrap(), SLOT_READY);
assert!(
slot.parked_thread.lock().is_none(),
"unparking a ready slot must clear the parked waiter handle"
);
}
#[test]
fn acquire_release_single_thread() {
let shard = FcPageLockShard::new(0);
let p = page(100_000);
let t = txn(1);
assert!(shard.try_acquire(p, t).is_ok());
assert_eq!(shard.holder(p), Some(t));
assert!(shard.release(p, t));
assert!(shard.holder(p).is_none());
}
#[test]
fn reacquire_by_holder_is_idempotent() {
let shard = FcPageLockShard::new(1);
let p = page(200_000);
let t = txn(2);
assert!(shard.try_acquire(p, t).is_ok());
assert!(shard.try_acquire(p, t).is_ok());
assert!(shard.release(p, t));
assert!(!shard.release(p, t));
}
#[test]
fn contending_txns_get_held_error() {
let shard = FcPageLockShard::new(2);
let p = page(300_000);
let holder = txn(1);
let other = txn(2);
assert!(shard.try_acquire(p, holder).is_ok());
assert_eq!(shard.try_acquire(p, other), Err(holder));
}
#[test]
fn release_all_drops_only_matching_txn() {
let shard = FcPageLockShard::new(3);
let t1 = txn(1);
let t2 = txn(2);
for i in 0..10u32 {
shard.try_acquire(page(400_000 + i), t1).unwrap();
}
for i in 0..5u32 {
shard.try_acquire(page(500_000 + i), t2).unwrap();
}
assert_eq!(shard.len(), 15);
let removed = shard.release_all(t1);
assert_eq!(removed, 10);
assert_eq!(shard.len(), 5);
}
#[test]
fn concurrent_acquire_distinct_pages() {
let shard = Arc::new(FcPageLockShard::new(4));
let barrier = Arc::new(Barrier::new(8));
let mut handles = Vec::new();
for t in 1..=8u64 {
let s = Arc::clone(&shard);
let b = Arc::clone(&barrier);
handles.push(thread::spawn(move || {
b.wait();
let tt = txn(t);
for i in 0..100u32 {
let p = page(600_000 + (t as u32 * 1000) + i);
s.try_acquire(p, tt).unwrap();
assert_eq!(s.holder(p), Some(tt));
assert!(s.release(p, tt));
}
}));
}
for h in handles {
h.join().unwrap();
}
assert_eq!(shard.len(), 0);
}
#[test]
fn concurrent_contention_single_page() {
let shard = Arc::new(FcPageLockShard::new(5));
let p = page(700_000);
let barrier = Arc::new(Barrier::new(8));
let mut handles = Vec::new();
for t in 1..=8u64 {
let s = Arc::clone(&shard);
let b = Arc::clone(&barrier);
handles.push(thread::spawn(move || {
b.wait();
let tt = txn(t);
let mut wins = 0u32;
for _ in 0..200u32 {
match s.try_acquire(p, tt) {
Ok(()) => {
wins += 1;
assert_eq!(s.holder(p), Some(tt));
assert!(s.release(p, tt));
}
Err(h) => assert_ne!(h, tt),
}
}
wins
}));
}
let total_wins: u32 = handles.into_iter().map(|h| h.join().unwrap()).sum();
assert!(total_wins > 0);
assert_eq!(shard.len(), 0);
}
#[test]
fn retain_drops_matching_entries() {
let shard = FcPageLockShard::new(6);
let t = txn(1);
for i in 0..20u32 {
shard.try_acquire(page(800_000 + i), t).unwrap();
}
let dropped = shard.retain(|p, _| p.get() % 2 == 0);
assert_eq!(dropped, 10);
assert_eq!(shard.len(), 10);
}
#[test]
fn holder_reports_unlocked_for_never_locked_page() {
let shard = FcPageLockShard::new(7);
assert!(shard.holder(page(1_000_000)).is_none());
}
}