use std::collections::HashMap;
use std::path::Path;
use parking_lot::Mutex;
use uqa_sql::ast::LockStrength;
use super::super::lock_strengths_conflict;
use super::{ByteClaim, PhysicalRowChangeTarget, RowChangeTarget};
use super::{PublishedRowChange, PublishedRowChangeKind, PublishedRowIdentity};
const CHANGE_JOURNAL_LOCK_BYTE: u64 = 10;
const SLOT_METADATA_LOCK_BYTE: u64 = 11;
#[cfg(windows)]
const MODE_TRANSITION_LOCK_BYTE: u64 = 12;
const TRANSACTION_XID_ATTACHMENT_BYTE: u64 = 8;
const TRANSACTION_XID_LOCK_BYTE: u64 = 13;
const TRANSACTION_XID_STATE_OFFSET: u64 = 16;
const TRANSACTION_XID_STATE_SIZE: usize = 16;
const TRANSACTION_XID_STATE_MAGIC: u32 = 0x5551_5849;
const TRANSACTION_XID_STATE_VERSION: u32 = 1;
const TRANSACTION_XID_CURSOR_OFFSET: u64 = 32;
const TRANSACTION_XID_CURSOR_SIZE: usize = 24;
const TRANSACTION_XID_CURSOR_MAGIC: u32 = 0x5551_5843;
const TRANSACTION_XID_CURSOR_VERSION: u32 = 1;
const WAIT_SLOT_BASE: u64 = 64;
const WAIT_SLOT_SIZE: u64 = 32;
const WAIT_SLOT_COUNT: u64 = 256;
const HOLDER_SLOT_BASE: u64 = WAIT_SLOT_BASE + WAIT_SLOT_SIZE * WAIT_SLOT_COUNT;
const HOLDER_SLOT_SIZE: u64 = 32;
const HOLDER_SLOT_COUNT: u64 = 8192;
const CHANGE_ENTRY_SIZE: u64 = 48;
const CHANGE_ENTRY_MAGIC: u32 = 0x5551_4348;
const CHANGE_JOURNAL_WAIT_LIMIT: std::time::Duration = std::time::Duration::from_secs(30);
#[derive(Default)]
struct ByteClaimCounts {
shared: u64,
exclusive: u64,
}
impl ByteClaimCounts {
fn mode(&self) -> Option<bool> {
if self.exclusive > 0 {
Some(true)
} else if self.shared > 0 {
Some(false)
} else {
None
}
}
}
struct CoordinatorState {
claims: HashMap<u64, ByteClaimCounts>,
holders: HashMap<u64, Vec<u64>>,
wait_slots: HashMap<u64, u64>,
holder_slots: HashMap<(u64, u64, bool), Vec<u64>>,
pending_holders: HashMap<(u64, u64, bool), Vec<u64>>,
pending_order: std::collections::BTreeMap<u64, (u64, ByteClaim)>,
next_pending: u64,
occupied_holder_slots: Vec<bool>,
released_holder_slots: Vec<u64>,
next_holder_slot: u64,
rows: row_claims::RowClaims,
}
pub(in crate::row_locks) struct FileLockCoordinator {
file: std::fs::File,
change_file: std::fs::File,
claim_file: std::fs::File,
sequence_file: std::fs::File,
change_journal: Mutex<()>,
transaction_xids: Mutex<xids::TransactionXids>,
temporary_role_slots: Mutex<temporary_roles::Slots>,
state: Mutex<CoordinatorState>,
}
mod claims;
mod journal;
mod platform;
mod relations;
mod row_claims;
mod sequence_positions;
mod temporary_roles;
mod waits;
mod xids;
use platform::{lock_would_block, process_alive, read_exact_at, write_all_at};
impl FileLockCoordinator {
pub(in crate::row_locks) fn open(database_path: &Path) -> Result<Self, String> {
let mut sidecar = database_path.as_os_str().to_owned();
sidecar.push(".uqa-locks");
let file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&sidecar)
.map_err(|error| {
format!(
"open cross-process lock file `{}`: {error}",
Path::new(&sidecar).display()
)
})?;
let mut change_sidecar = database_path.as_os_str().to_owned();
change_sidecar.push(".uqa-row-changes");
let change_file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&change_sidecar)
.map_err(|error| {
format!(
"open cross-process row-change journal `{}`: {error}",
Path::new(&change_sidecar).display()
)
})?;
let mut claim_sidecar = database_path.as_os_str().to_owned();
claim_sidecar.push(".uqa-row-claims");
let claim_file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&claim_sidecar)
.map_err(|error| {
format!(
"open cross-process row claim table `{}`: {error}",
Path::new(&claim_sidecar).display()
)
})?;
let mut sequence_sidecar = database_path.as_os_str().to_owned();
sequence_sidecar.push(".uqa-sequences");
let sequence_file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&sequence_sidecar)
.map_err(|error| {
format!(
"open cross-process sequence positions `{}`: {error}",
Path::new(&sequence_sidecar).display()
)
})?;
let pid = std::process::id();
let coordinator = Self {
file,
change_file,
claim_file,
sequence_file,
change_journal: Mutex::new(()),
transaction_xids: Mutex::new(xids::TransactionXids::new()),
temporary_role_slots: Mutex::new(temporary_roles::Slots::default()),
state: Mutex::new(CoordinatorState {
claims: HashMap::new(),
holders: HashMap::new(),
wait_slots: HashMap::new(),
holder_slots: HashMap::new(),
pending_holders: HashMap::new(),
pending_order: std::collections::BTreeMap::new(),
next_pending: 0,
occupied_holder_slots: vec![false; HOLDER_SLOT_COUNT as usize],
released_holder_slots: Vec::new(),
next_holder_slot: u64::from(pid).wrapping_mul(31) % HOLDER_SLOT_COUNT,
rows: row_claims::RowClaims::default(),
}),
};
Ok(coordinator)
}
}
impl FileLockCoordinator {
#[allow(clippy::unused_self)]
pub(in crate::row_locks) const fn shares_sequence_positions(&self) -> bool {
true
}
}
impl Drop for FileLockCoordinator {
fn drop(&mut self) {
self.detach_transaction_xids();
self.detach_row_claims_process();
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn holder_slot_cursor_tracks_bulk_claims_and_releases() {
let directory = tempfile::tempdir().unwrap();
let coordinator = FileLockCoordinator::open(&directory.path().join("bulk.db")).unwrap();
let claims = (0..128)
.map(|ordinal| ByteClaim {
offset: 10_000 + ordinal,
write: true,
})
.collect::<Vec<_>>();
let session = 17;
let mut state = coordinator.state.lock();
let first_slot = state.next_holder_slot;
for claim in &claims {
coordinator.register_holder_slot(&mut state, session, *claim);
}
assert_eq!(
state
.occupied_holder_slots
.iter()
.filter(|occupied| **occupied)
.count(),
claims.len()
);
for (ordinal, claim) in claims.iter().enumerate() {
let expected = (first_slot + ordinal as u64) % HOLDER_SLOT_COUNT;
assert_eq!(
state
.holder_slots
.get(&(session, claim.offset, claim.write))
.unwrap(),
&[expected]
);
}
assert_eq!(
state.next_holder_slot,
(first_slot + claims.len() as u64) % HOLDER_SLOT_COUNT
);
for claim in &claims {
coordinator.clear_holder_slots(&mut state, session, std::slice::from_ref(claim));
}
assert!(state.holder_slots.is_empty());
assert!(state.occupied_holder_slots.iter().all(|occupied| !occupied));
}
#[test]
fn released_holder_slots_bound_file_growth_across_transactions() {
let directory = tempfile::tempdir().unwrap();
let coordinator = FileLockCoordinator::open(&directory.path().join("reuse.db")).unwrap();
coordinator.state.lock().next_holder_slot = 0;
for transaction in 0..4 {
let claims = (0..128)
.map(|ordinal| ByteClaim {
offset: 10_000 + transaction * 128 + ordinal,
write: true,
})
.collect::<Vec<_>>();
assert!(matches!(coordinator.try_claim(17, &claims), Ok(Ok(()))));
coordinator.publish_holders();
assert_eq!(
coordinator.file.metadata().unwrap().len(),
HOLDER_SLOT_BASE + HOLDER_SLOT_SIZE * 128
);
coordinator.release(17, &claims);
let state = coordinator.state.lock();
assert!(state.holder_slots.is_empty());
assert!(state.claims.is_empty());
assert_eq!(state.next_holder_slot, 128);
assert_eq!(state.released_holder_slots.len(), 128);
}
}
#[test]
fn holders_are_published_before_any_local_wait_is_advertised() {
let directory = tempfile::tempdir().unwrap();
let coordinator = FileLockCoordinator::open(&directory.path().join("pending.db")).unwrap();
let held = [10_000, 10_001].map(|offset| ByteClaim {
offset,
write: true,
});
let released = ByteClaim {
offset: 10_002,
write: false,
};
assert!(matches!(coordinator.try_claim(17, &held), Ok(Ok(()))));
assert!(matches!(coordinator.try_claim(17, &[released]), Ok(Ok(()))));
coordinator.release(17, &[released]);
{
let state = coordinator.state.lock();
assert!(state.holder_slots.is_empty(), "claims write no slot");
assert!(
state.released_holder_slots.is_empty(),
"an unpublished release clears no slot"
);
}
coordinator.register_wait(
23,
ByteClaim {
offset: 20_000,
write: true,
},
);
let state = coordinator.state.lock();
assert!(state.pending_holders.is_empty() && state.pending_order.is_empty());
for claim in held {
let index = state.holder_slots[&(17, claim.offset, claim.write)][0];
let slot = coordinator.read_holder_slot(index).unwrap();
assert_eq!((slot.session, slot.offset), (17, claim.offset));
}
assert!(!state
.holder_slots
.contains_key(&(17, released.offset, released.write)));
}
#[test]
fn released_holes_are_reused_without_overwriting_live_holders() {
let directory = tempfile::tempdir().unwrap();
let coordinator = FileLockCoordinator::open(&directory.path().join("holes.db")).unwrap();
coordinator.state.lock().next_holder_slot = 0;
let first = (0..32)
.map(|ordinal| ByteClaim {
offset: 10_000 + ordinal,
write: true,
})
.collect::<Vec<_>>();
assert!(matches!(coordinator.try_claim(17, &first), Ok(Ok(()))));
coordinator.publish_holders();
for claim in first.iter().step_by(2) {
coordinator.release(17, &[*claim]);
}
let second = (0..16)
.map(|ordinal| ByteClaim {
offset: 20_000 + ordinal,
write: true,
})
.collect::<Vec<_>>();
assert!(matches!(coordinator.try_claim(23, &second), Ok(Ok(()))));
coordinator.publish_holders();
{
let state = coordinator.state.lock();
for (ordinal, claim) in first.iter().enumerate().skip(1).step_by(2) {
assert_eq!(
state.holder_slots[&(17, claim.offset, claim.write)],
[ordinal as u64]
);
}
let reused = second
.iter()
.map(|claim| state.holder_slots[&(23, claim.offset, claim.write)][0])
.collect::<std::collections::BTreeSet<_>>();
assert_eq!(reused, (0..32).step_by(2).collect());
assert_eq!(state.next_holder_slot, 32);
}
assert_eq!(
coordinator.file.metadata().unwrap().len(),
HOLDER_SLOT_BASE + HOLDER_SLOT_SIZE * 32
);
for claim in first.iter().skip(1).step_by(2) {
coordinator.release(17, &[*claim]);
}
coordinator.release(23, &second);
assert!(coordinator.state.lock().holder_slots.is_empty());
}
}