use crate::db::WriteStrategy;
use crate::tree_store::btree_base::Checksum;
use crate::tree_store::page_store::buddy_allocator::BuddyAllocator;
use crate::tree_store::page_store::layout::{DatabaseLayout, RegionLayout};
use crate::tree_store::page_store::mmap::Mmap;
use crate::tree_store::page_store::page_allocator::PageAllocator;
use crate::tree_store::page_store::region::{RegionHeaderAccessor, RegionHeaderMutator};
use crate::tree_store::page_store::utils::get_page_size;
use crate::tree_store::page_store::{hash128_with_seed, PageImpl, PageMut};
use crate::tree_store::PageNumber;
use crate::Error;
use crate::Result;
use std::cmp::{max, min};
#[cfg(debug_assertions)]
use std::collections::HashMap;
use std::collections::HashSet;
use std::convert::TryInto;
use std::fs::File;
use std::io;
use std::io::{Read, Seek, SeekFrom};
use std::mem::size_of;
use std::path::Path;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Mutex, MutexGuard};
const MAX_USABLE_REGION_SPACE: u64 = 4 * 1024 * 1024 * 1024;
pub(crate) const MAX_MAX_PAGE_ORDER: usize = 20;
pub(super) const MIN_USABLE_PAGES: usize = 10;
const MIN_DESIRED_USABLE_BYTES: usize = 1024 * 1024;
const FILE_FORMAT_VERSION: u8 = 105;
const MAGICNUMBER: [u8; 9] = [b'r', b'e', b'd', b'b', 0x1A, 0x0A, 0xA9, 0x0D, 0x0A];
const GOD_BYTE_OFFSET: usize = MAGICNUMBER.len();
const CHECKSUM_TYPE_OFFSET: usize = GOD_BYTE_OFFSET + size_of::<u8>();
const SUPERHEADER_PAGES_OFFSET: usize = CHECKSUM_TYPE_OFFSET + size_of::<u8>() + 1; const PAGE_SIZE_OFFSET: usize = SUPERHEADER_PAGES_OFFSET + size_of::<u32>();
const REGION_TRACKER_LENGTH_OFFSET: usize = PAGE_SIZE_OFFSET + size_of::<u32>();
const DB_SIZE_OFFSET: usize = REGION_TRACKER_LENGTH_OFFSET + size_of::<u32>();
const REGION_HEADER_PAGES_OFFSET: usize = DB_SIZE_OFFSET + size_of::<u64>();
const REGION_MAX_DATA_PAGES_OFFSET: usize = REGION_HEADER_PAGES_OFFSET + size_of::<u32>();
const TRANSACTION_SIZE: usize = 128;
const TRANSACTION_0_OFFSET: usize = 64;
const TRANSACTION_1_OFFSET: usize = TRANSACTION_0_OFFSET + TRANSACTION_SIZE;
pub(super) const DB_HEADER_SIZE: usize = TRANSACTION_1_OFFSET + TRANSACTION_SIZE;
const PRIMARY_BIT: u8 = 1;
const RECOVERY_REQUIRED: u8 = 2;
const VERSION_OFFSET: usize = 0;
const ROOT_NON_NULL_OFFSET: usize = size_of::<u8>();
const FREED_ROOT_NON_NULL_OFFSET: usize = ROOT_NON_NULL_OFFSET + size_of::<u8>();
const PADDING: usize = 5;
const ROOT_PAGE_OFFSET: usize = FREED_ROOT_NON_NULL_OFFSET + size_of::<u8>() + PADDING;
const ROOT_CHECKSUM_OFFSET: usize = ROOT_PAGE_OFFSET + size_of::<u64>();
const FREED_ROOT_OFFSET: usize = ROOT_CHECKSUM_OFFSET + size_of::<u128>();
const FREED_ROOT_CHECKSUM_OFFSET: usize = FREED_ROOT_OFFSET + size_of::<u64>();
const TRANSACTION_ID_OFFSET: usize = FREED_ROOT_CHECKSUM_OFFSET + size_of::<u128>();
const NUM_FULL_REGIONS_OFFSET: usize = TRANSACTION_ID_OFFSET + size_of::<u64>();
const TRAILING_REGION_DATA_PAGES_OFFSET: usize = NUM_FULL_REGIONS_OFFSET + size_of::<u32>();
const SLOT_CHECKSUM_OFFSET: usize = TRAILING_REGION_DATA_PAGES_OFFSET + size_of::<u32>();
const TRANSACTION_LAST_FIELD: usize = SLOT_CHECKSUM_OFFSET + size_of::<u128>();
fn ceil_log2(x: usize) -> usize {
if x.is_power_of_two() {
x.trailing_zeros() as usize
} else {
x.next_power_of_two().trailing_zeros() as usize
}
}
pub(crate) fn get_db_size(path: impl AsRef<Path>) -> Result<u64, io::Error> {
let mut db_size = [0u8; size_of::<u64>()];
let mut file = File::open(path)?;
file.seek(SeekFrom::Start(DB_SIZE_OFFSET as u64))?;
file.read_exact(&mut db_size)?;
Ok(u64::from_le_bytes(db_size))
}
#[derive(Debug, PartialEq, Copy, Clone)]
pub(crate) enum ChecksumType {
Zero, XXH3_128,
}
impl ChecksumType {
pub(crate) fn checksum(&self, data: &[u8]) -> Checksum {
match self {
ChecksumType::Zero => 0,
ChecksumType::XXH3_128 => hash128_with_seed(data, 0),
}
}
}
impl From<u8> for ChecksumType {
fn from(x: u8) -> Self {
match x {
1 => ChecksumType::Zero,
2 => ChecksumType::XXH3_128,
_ => unimplemented!(),
}
}
}
#[allow(clippy::from_over_into)]
impl Into<u8> for ChecksumType {
fn into(self) -> u8 {
match self {
ChecksumType::Zero => 1,
ChecksumType::XXH3_128 => 2,
}
}
}
struct MetadataGuard;
struct MetadataAccessor<'a> {
header: &'a mut [u8],
mmap: &'a Mmap,
guard: MutexGuard<'a, MetadataGuard>,
}
impl<'a> MetadataAccessor<'a> {
unsafe fn new(mmap: &'a Mmap, guard: MutexGuard<'a, MetadataGuard>) -> Self {
let header = mmap.get_memory_mut(0..DB_HEADER_SIZE);
Self {
header,
mmap,
guard,
}
}
fn make_region_layout(&self, num_pages: usize) -> RegionLayout {
RegionLayout::new(
num_pages,
self.get_region_max_data_pages(),
self.get_region_header_length(),
self.get_page_size(),
)
}
fn get_primary_layout(&self) -> DatabaseLayout {
let trailing_pages = self.primary_slot().get_trailing_region_data_pages();
let num_full_regions = self.primary_slot().get_full_regions();
let full_region = self.make_region_layout(self.get_region_max_data_pages());
let trailing_region = trailing_pages.map(|x| self.make_region_layout(x));
DatabaseLayout::new(
self.get_superheader_length(),
self.get_region_tracker_state_length(),
num_full_regions,
full_region,
trailing_region,
)
}
fn get_secondary_layout(&self) -> DatabaseLayout {
let trailing_pages = self.secondary_slot().get_trailing_region_data_pages();
let num_full_regions = self.secondary_slot().get_full_regions();
let full_region = self.make_region_layout(self.get_region_max_data_pages());
let trailing_region = trailing_pages.map(|x| self.make_region_layout(x));
DatabaseLayout::new(
self.get_superheader_length(),
self.get_region_tracker_state_length(),
num_full_regions,
full_region,
trailing_region,
)
}
fn primary_slot(&self) -> TransactionAccessor {
let start = if self.header[GOD_BYTE_OFFSET] & PRIMARY_BIT == 0 {
TRANSACTION_0_OFFSET
} else {
TRANSACTION_1_OFFSET
};
let end = start + TRANSACTION_SIZE;
let mem = &self.header[start..end];
TransactionAccessor::new(mem, &self.guard)
}
fn secondary_slot(&self) -> TransactionAccessor {
let start = if self.header[GOD_BYTE_OFFSET] & PRIMARY_BIT == 0 {
TRANSACTION_1_OFFSET
} else {
TRANSACTION_0_OFFSET
};
let end = start + TRANSACTION_SIZE;
let mem = &self.header[start..end];
TransactionAccessor::new(mem, &self.guard)
}
fn secondary_slot_mut(&mut self) -> TransactionMutator {
let start = if self.header[GOD_BYTE_OFFSET] & PRIMARY_BIT == 0 {
TRANSACTION_1_OFFSET
} else {
TRANSACTION_0_OFFSET
};
let end = start + TRANSACTION_SIZE;
let mem = &mut self.header[start..end];
TransactionMutator::new(mem)
}
fn swap_primary(&mut self) {
if self.header[GOD_BYTE_OFFSET] & PRIMARY_BIT == 0 {
self.header[GOD_BYTE_OFFSET] |= PRIMARY_BIT;
} else {
self.header[GOD_BYTE_OFFSET] &= !PRIMARY_BIT;
}
}
fn get_checksum_type(&self) -> ChecksumType {
ChecksumType::from(self.header[CHECKSUM_TYPE_OFFSET])
}
fn set_checksum_type(&mut self, checksum: ChecksumType) {
self.header[CHECKSUM_TYPE_OFFSET] = checksum.into();
}
fn get_max_capacity(&self) -> u64 {
u64::from_le_bytes(
self.header[DB_SIZE_OFFSET..DB_SIZE_OFFSET + size_of::<u64>()]
.try_into()
.unwrap(),
)
}
fn set_max_capacity(&mut self, max_size: u64) {
self.header[DB_SIZE_OFFSET..DB_SIZE_OFFSET + size_of::<u64>()]
.copy_from_slice(&max_size.to_le_bytes());
}
fn get_magic_number(&self) -> [u8; MAGICNUMBER.len()] {
self.header[..MAGICNUMBER.len()].try_into().unwrap()
}
fn set_magic_number(&mut self) {
self.header[..MAGICNUMBER.len()].copy_from_slice(&MAGICNUMBER);
}
fn get_superheader_length(&self) -> usize {
let pages = u32::from_le_bytes(
self.header[SUPERHEADER_PAGES_OFFSET..(SUPERHEADER_PAGES_OFFSET + size_of::<u32>())]
.try_into()
.unwrap(),
) as usize;
pages * self.get_page_size()
}
fn set_superheader_length(&mut self, length: usize) {
assert_eq!(length % self.get_page_size(), 0);
let pages = length / self.get_page_size();
self.header[SUPERHEADER_PAGES_OFFSET..(SUPERHEADER_PAGES_OFFSET + size_of::<u32>())]
.copy_from_slice(&(pages as u32).to_le_bytes());
}
fn get_region_tracker_state_length(&self) -> usize {
u32::from_le_bytes(
self.header
[REGION_TRACKER_LENGTH_OFFSET..(REGION_TRACKER_LENGTH_OFFSET + size_of::<u32>())]
.try_into()
.unwrap(),
) as usize
}
fn set_region_tracker_state_length(&mut self, length: usize) {
self.header
[REGION_TRACKER_LENGTH_OFFSET..(REGION_TRACKER_LENGTH_OFFSET + size_of::<u32>())]
.copy_from_slice(&(length as u32).to_le_bytes());
}
fn get_region_header_length(&self) -> usize {
let pages = u32::from_le_bytes(
self.header
[REGION_HEADER_PAGES_OFFSET..(REGION_HEADER_PAGES_OFFSET + size_of::<u32>())]
.try_into()
.unwrap(),
) as usize;
pages * self.get_page_size()
}
fn set_region_header_length(&mut self, length: usize) {
assert_eq!(length % self.get_page_size(), 0);
let pages = length / self.get_page_size();
self.header[REGION_HEADER_PAGES_OFFSET..(REGION_HEADER_PAGES_OFFSET + size_of::<u32>())]
.copy_from_slice(&(pages as u32).to_le_bytes());
}
fn get_region_max_usable_bytes(&self) -> u64 {
self.get_region_max_data_pages() as u64 * self.get_page_size() as u64
}
fn get_region_max_data_pages(&self) -> usize {
u32::from_le_bytes(
self.header
[REGION_MAX_DATA_PAGES_OFFSET..(REGION_MAX_DATA_PAGES_OFFSET + size_of::<u32>())]
.try_into()
.unwrap(),
) as usize
}
fn set_region_max_data_pages(&mut self, pages: usize) {
self.header
[REGION_MAX_DATA_PAGES_OFFSET..(REGION_MAX_DATA_PAGES_OFFSET + size_of::<u32>())]
.copy_from_slice(&(pages as u32).to_le_bytes());
}
fn get_page_size(&self) -> usize {
u32::from_le_bytes(
self.header[PAGE_SIZE_OFFSET..(PAGE_SIZE_OFFSET + size_of::<u32>())]
.try_into()
.unwrap(),
) as usize
}
fn set_page_size(&mut self, page_size: usize) {
self.header[PAGE_SIZE_OFFSET..(PAGE_SIZE_OFFSET + size_of::<u32>())]
.copy_from_slice(&(page_size as u32).to_le_bytes());
}
fn get_recovery_required(&self) -> bool {
self.header[GOD_BYTE_OFFSET] & RECOVERY_REQUIRED != 0
}
fn set_recovery(&mut self, required: bool) {
if required {
self.header[GOD_BYTE_OFFSET] |= RECOVERY_REQUIRED;
} else {
self.header[GOD_BYTE_OFFSET] &= !RECOVERY_REQUIRED;
}
}
fn get_region(&mut self, region: usize, layout: &DatabaseLayout) -> RegionHeaderAccessor {
let base = layout.region_base_address(region);
let len = layout.region_layout(region).data_section().start;
let absolute = base..(base + len);
let mem = unsafe { self.mmap.get_memory(absolute) };
RegionHeaderAccessor::new(mem)
}
fn initialize_region_tracker(&mut self, layout: &DatabaseLayout) {
let max_regions = DatabaseLayout::calculate(
self.get_max_capacity(),
self.get_max_capacity(),
self.get_region_max_usable_bytes(),
self.get_page_size(),
)
.unwrap()
.num_regions();
let range = layout.region_tracker_address_range();
assert!(range.start >= DB_HEADER_SIZE);
let mem = unsafe { self.mmap.get_memory_mut(range) };
RegionTracker::init_new(max_regions, MAX_MAX_PAGE_ORDER + 1, mem);
}
fn allocators_mut(
&mut self,
layout: &DatabaseLayout,
) -> Result<(RegionTracker, RegionsAccessor)> {
if !self.get_recovery_required() {
self.set_recovery(true);
self.mmap.flush()?
}
let range = layout.region_tracker_address_range();
assert!(range.start >= DB_HEADER_SIZE);
let mem = unsafe { self.mmap.get_memory_mut(range) };
let region_accessor = RegionsAccessor {
mmap: self.mmap,
layout: layout.clone(),
};
Ok((RegionTracker::new(mem), region_accessor))
}
}
pub(crate) struct RegionTracker<'a> {
accessor: PageAllocator,
data: &'a mut [u8],
}
impl<'a> RegionTracker<'a> {
pub(crate) fn new(data: &'a mut [u8]) -> Self {
let regions = u32::from_le_bytes(data[8..12].try_into().unwrap()) as usize;
Self {
accessor: PageAllocator::new(regions),
data,
}
}
pub(crate) fn required_bytes(regions: usize, orders: usize) -> usize {
2 * size_of::<u32>() + orders * PageAllocator::required_space(regions)
}
pub(crate) fn init_new(regions: usize, orders: usize, data: &'a mut [u8]) -> Self {
assert!(data.len() >= Self::required_bytes(regions, orders));
data[..4].copy_from_slice(&(orders as u32).to_le_bytes());
data[4..8].copy_from_slice(&(PageAllocator::required_space(regions) as u32).to_le_bytes());
let mut result = Self {
accessor: PageAllocator::new(regions),
data,
};
for i in 0..orders {
PageAllocator::init_new(result.get_order_mut(i), regions);
}
result
}
pub(crate) fn find_free(&self, order: usize) -> Result<u64> {
self.accessor.find_free(self.get_order(order))
}
pub(crate) fn mark_free(&mut self, order: usize, region: u64) {
assert!(order < self.suballocators());
for i in 0..=order {
let start = 8 + i * self.suballocator_len();
let end = start + self.suballocator_len();
let mem = &mut self.data[start..end];
self.accessor.free(mem, region);
}
}
pub(crate) fn mark_full(&mut self, order: usize, region: u64) {
assert!(order < self.suballocators());
for i in order..self.suballocators() {
let start = 8 + i * self.suballocator_len();
let end = start + self.suballocator_len();
let mem = &mut self.data[start..end];
self.accessor.record_alloc(mem, region);
}
}
fn suballocator_len(&self) -> usize {
u32::from_le_bytes(self.data[4..8].try_into().unwrap()) as usize
}
fn suballocators(&self) -> usize {
u32::from_le_bytes(self.data[..4].try_into().unwrap()) as usize
}
fn get_order_mut(&mut self, order: usize) -> &mut [u8] {
assert!(order < self.suballocators());
let start = 8 + order * self.suballocator_len();
let end = start + self.suballocator_len();
&mut self.data[start..end]
}
fn get_order(&self, order: usize) -> &[u8] {
assert!(order < self.suballocators());
let start = 8 + order * self.suballocator_len();
let end = start + self.suballocator_len();
&self.data[start..end]
}
}
struct RegionsAccessor<'a> {
mmap: &'a Mmap,
layout: DatabaseLayout,
}
impl<'a> RegionsAccessor<'a> {
fn get_region_mut(&mut self, region: usize) -> RegionHeaderMutator {
let base = self.layout.region_base_address(region);
let region_header_len = &self.layout.region_layout(region).data_section().start;
let absolute = base..(base + region_header_len);
assert!(absolute.start >= self.layout.header_bytes());
let mem = unsafe { self.mmap.get_memory_mut(absolute) };
RegionHeaderMutator::new(mem)
}
}
struct TransactionAccessor<'a> {
mem: &'a [u8],
_guard: &'a MutexGuard<'a, MetadataGuard>,
}
impl<'a> TransactionAccessor<'a> {
fn new(mem: &'a [u8], guard: &'a MutexGuard<'a, MetadataGuard>) -> Self {
TransactionAccessor { mem, _guard: guard }
}
fn verify_checksum(&self, checksum_type: ChecksumType) -> bool {
let checksum = Checksum::from_le_bytes(
self.mem[SLOT_CHECKSUM_OFFSET..(SLOT_CHECKSUM_OFFSET + size_of::<Checksum>())]
.try_into()
.unwrap(),
);
checksum_type.checksum(&self.mem[..SLOT_CHECKSUM_OFFSET]) == checksum
}
fn get_root_page(&self) -> Option<(PageNumber, Checksum)> {
if self.mem[ROOT_NON_NULL_OFFSET] == 0 {
None
} else {
let num = PageNumber::from_le_bytes(
self.mem[ROOT_PAGE_OFFSET..(ROOT_PAGE_OFFSET + PageNumber::serialized_size())]
.try_into()
.unwrap(),
);
let checksum = Checksum::from_le_bytes(
self.mem[ROOT_CHECKSUM_OFFSET..(ROOT_CHECKSUM_OFFSET + size_of::<Checksum>())]
.try_into()
.unwrap(),
);
Some((num, checksum))
}
}
fn get_freed_root_page(&self) -> Option<(PageNumber, Checksum)> {
if self.mem[FREED_ROOT_NON_NULL_OFFSET] == 0 {
None
} else {
let num = PageNumber::from_le_bytes(
self.mem[FREED_ROOT_OFFSET..(FREED_ROOT_OFFSET + PageNumber::serialized_size())]
.try_into()
.unwrap(),
);
let checksum = Checksum::from_le_bytes(
self.mem[FREED_ROOT_CHECKSUM_OFFSET
..(FREED_ROOT_CHECKSUM_OFFSET + size_of::<Checksum>())]
.try_into()
.unwrap(),
);
Some((num, checksum))
}
}
fn get_last_committed_transaction_id(&self) -> u64 {
u64::from_le_bytes(
self.mem[TRANSACTION_ID_OFFSET..(TRANSACTION_ID_OFFSET + size_of::<u64>())]
.try_into()
.unwrap(),
)
}
fn get_full_regions(&self) -> usize {
u32::from_le_bytes(
self.mem[NUM_FULL_REGIONS_OFFSET..(NUM_FULL_REGIONS_OFFSET + size_of::<u32>())]
.try_into()
.unwrap(),
) as usize
}
fn get_trailing_region_data_pages(&self) -> Option<usize> {
let value = u32::from_le_bytes(
self.mem[TRAILING_REGION_DATA_PAGES_OFFSET
..(TRAILING_REGION_DATA_PAGES_OFFSET + size_of::<u32>())]
.try_into()
.unwrap(),
);
if value == 0 {
None
} else {
Some(value as usize)
}
}
fn get_version(&self) -> u8 {
self.mem[VERSION_OFFSET]
}
}
struct TransactionMutator<'a> {
mem: &'a mut [u8],
}
impl<'a> TransactionMutator<'a> {
fn new(mem: &'a mut [u8]) -> Self {
TransactionMutator { mem }
}
fn set_root_page(&mut self, page_number: Option<(PageNumber, Checksum)>) {
if let Some((num, checksum)) = page_number {
self.mem[ROOT_PAGE_OFFSET..(ROOT_PAGE_OFFSET + PageNumber::serialized_size())]
.copy_from_slice(&num.to_le_bytes());
self.mem[ROOT_CHECKSUM_OFFSET..(ROOT_CHECKSUM_OFFSET + size_of::<Checksum>())]
.copy_from_slice(&checksum.to_le_bytes());
self.mem[ROOT_NON_NULL_OFFSET] = 1;
} else {
self.mem[ROOT_NON_NULL_OFFSET] = 0;
}
}
fn set_freed_root(&mut self, page_number: Option<(PageNumber, Checksum)>) {
if let Some((num, checksum)) = page_number {
self.mem[FREED_ROOT_OFFSET..(FREED_ROOT_OFFSET + PageNumber::serialized_size())]
.copy_from_slice(&num.to_le_bytes());
self.mem
[FREED_ROOT_CHECKSUM_OFFSET..(FREED_ROOT_CHECKSUM_OFFSET + size_of::<Checksum>())]
.copy_from_slice(&checksum.to_le_bytes());
self.mem[FREED_ROOT_NON_NULL_OFFSET] = 1;
} else {
self.mem[FREED_ROOT_NON_NULL_OFFSET] = 0;
}
}
fn set_last_committed_transaction_id(&mut self, transaction_id: u64) {
self.mem[TRANSACTION_ID_OFFSET..(TRANSACTION_ID_OFFSET + size_of::<u64>())]
.copy_from_slice(&transaction_id.to_le_bytes());
}
fn set_data_section_layout(
&mut self,
full_regions: usize,
trailing_region_data_pages: Option<usize>,
) {
self.mem[NUM_FULL_REGIONS_OFFSET..(NUM_FULL_REGIONS_OFFSET + size_of::<u32>())]
.copy_from_slice(&(full_regions as u32).to_le_bytes());
self.mem[TRAILING_REGION_DATA_PAGES_OFFSET
..(TRAILING_REGION_DATA_PAGES_OFFSET + size_of::<u32>())]
.copy_from_slice(&(trailing_region_data_pages.unwrap_or(0) as u32).to_le_bytes());
}
fn update_checksum(&mut self, checksum_type: ChecksumType) {
let checksum = checksum_type.checksum(&self.mem[..SLOT_CHECKSUM_OFFSET]);
self.mem[SLOT_CHECKSUM_OFFSET..(SLOT_CHECKSUM_OFFSET + size_of::<Checksum>())]
.copy_from_slice(&checksum.to_le_bytes());
}
fn set_version(&mut self, version: u8) {
self.mem[VERSION_OFFSET] = version;
}
}
enum AllocationOp {
Allocate(PageNumber),
Free(PageNumber),
FreeUncommitted(PageNumber),
}
pub(crate) struct TransactionalMemory {
allocated_since_commit: Mutex<HashSet<PageNumber>>,
log_since_commit: Mutex<Vec<AllocationOp>>,
regional_allocators: Mutex<Option<Vec<BuddyAllocator>>>,
mmap: Mmap,
metadata_guard: Mutex<MetadataGuard>,
layout: Mutex<DatabaseLayout>,
#[cfg(debug_assertions)]
open_dirty_pages: Mutex<HashSet<PageNumber>>,
#[cfg(debug_assertions)]
read_page_ref_counts: Mutex<HashMap<PageNumber, u64>>,
read_from_secondary: AtomicBool,
page_size: usize,
region_size: u64,
region_header_with_padding_size: usize,
db_header_size: usize,
dynamic_growth: bool,
checksum_type: ChecksumType,
}
impl TransactionalMemory {
pub(crate) fn new(
file: File,
max_capacity: u64,
requested_page_size: Option<usize>,
requested_region_size: Option<usize>,
dynamic_growth: bool,
write_strategy: Option<WriteStrategy>,
) -> Result<Self> {
#[allow(clippy::assertions_on_constants)]
{
assert!(TRANSACTION_LAST_FIELD <= TRANSACTION_SIZE);
}
let page_size = requested_page_size.unwrap_or_else(get_page_size);
assert!(page_size.is_power_of_two());
if max_capacity < (DB_HEADER_SIZE + page_size * MIN_USABLE_PAGES) as u64 {
return Err(Error::OutOfSpace);
}
let mmap = Mmap::new(file, max_capacity.try_into().unwrap())?;
if mmap.len() < DB_HEADER_SIZE {
unsafe {
mmap.resize(DB_HEADER_SIZE)?;
}
}
let mutex = Mutex::new(MetadataGuard {});
let mut metadata = unsafe { MetadataAccessor::new(&mmap, mutex.lock().unwrap()) };
let region_size = requested_region_size
.map(|x| x as u64)
.unwrap_or(MAX_USABLE_REGION_SPACE);
assert!(region_size.is_power_of_two());
let max_usable_region_bytes = min(region_size, max_capacity.next_power_of_two() as u64);
if metadata.get_magic_number() != MAGICNUMBER {
let starting_size = if dynamic_growth {
MIN_DESIRED_USABLE_BYTES as u64
} else {
max_capacity
};
let layout = DatabaseLayout::calculate(
max_capacity,
starting_size,
max_usable_region_bytes,
page_size,
)?;
if (mmap.len() as u64) < layout.len() {
unsafe {
mmap.resize(layout.len().try_into().unwrap())?;
}
}
metadata.header.fill(0);
metadata.set_page_size(page_size);
metadata.set_superheader_length(layout.header_bytes());
metadata.set_region_tracker_state_length(layout.region_tracker_address_range().len());
metadata.set_region_header_length(layout.full_region_layout().data_section().start);
metadata.set_region_max_data_pages(layout.full_region_layout().num_pages());
metadata.set_max_capacity(max_capacity);
let checksum_type = match write_strategy.unwrap_or_default() {
WriteStrategy::Checksum => ChecksumType::XXH3_128,
WriteStrategy::TwoPhase => ChecksumType::Zero,
};
metadata.set_checksum_type(checksum_type);
metadata.initialize_region_tracker(&layout);
let (mut region_tracker, mut regions) = metadata.allocators_mut(&layout)?;
let num_regions = layout.num_regions();
for i in 0..num_regions {
region_tracker.mark_free(layout.full_region_layout().max_order(), i as u64);
}
for i in 0..num_regions {
let mut region = regions.get_region_mut(i);
let region_layout = layout.region_layout(i);
region.initialize(
region_layout.num_pages(),
layout.full_region_layout().num_pages(),
region_layout.max_order(),
);
}
metadata.set_recovery(false);
let mut mutator = metadata.secondary_slot_mut();
mutator.set_root_page(None);
mutator.set_freed_root(None);
mutator.set_last_committed_transaction_id(0);
mutator.set_data_section_layout(
layout.num_full_regions(),
layout.trailing_region_layout().map(|x| x.num_pages()),
);
mutator.set_version(FILE_FORMAT_VERSION);
drop(mutator);
metadata.swap_primary();
let mut mutator = metadata.secondary_slot_mut();
mutator.set_data_section_layout(
layout.num_full_regions(),
layout.trailing_region_layout().map(|x| x.num_pages()),
);
mutator.set_version(FILE_FORMAT_VERSION);
drop(mutator);
mmap.flush()?;
metadata.set_magic_number();
mmap.flush()?;
}
let page_size = metadata.get_page_size();
if let Some(size) = requested_page_size {
assert_eq!(page_size, size);
}
let version = metadata.primary_slot().get_version();
if version != FILE_FORMAT_VERSION {
return Err(Error::Corrupted(format!(
"Expected file format version {}, found {}",
FILE_FORMAT_VERSION, version
)));
}
let version = metadata.secondary_slot().get_version();
if version != FILE_FORMAT_VERSION {
return Err(Error::Corrupted(format!(
"Expected file format version {}, found {}",
FILE_FORMAT_VERSION, version
)));
}
let layout = metadata.get_primary_layout();
let region_size = layout.full_region_layout().len();
let region_header_size = layout.full_region_layout().data_section().start;
let checksum_type = metadata.get_checksum_type();
let regional_allocators = if metadata.get_recovery_required() {
None
} else {
Some(layout.create_allocators())
};
drop(metadata);
Ok(TransactionalMemory {
allocated_since_commit: Mutex::new(HashSet::new()),
log_since_commit: Mutex::new(vec![]),
regional_allocators: Mutex::new(regional_allocators),
mmap,
metadata_guard: mutex,
layout: Mutex::new(layout.clone()),
#[cfg(debug_assertions)]
open_dirty_pages: Mutex::new(HashSet::new()),
#[cfg(debug_assertions)]
read_page_ref_counts: Mutex::new(HashMap::new()),
read_from_secondary: AtomicBool::new(false),
page_size,
region_size,
region_header_with_padding_size: region_header_size,
db_header_size: layout.header_bytes(),
dynamic_growth,
checksum_type,
})
}
pub(crate) fn needs_repair(&self) -> Result<bool> {
Ok(self.lock_metadata().get_recovery_required())
}
pub(crate) fn needs_checksum_verification(&self) -> Result<bool> {
Ok(self.lock_metadata().get_checksum_type() == ChecksumType::XXH3_128)
}
pub(crate) fn checksum_type(&self) -> ChecksumType {
self.lock_metadata().get_checksum_type()
}
pub(crate) fn repair_primary_corrupted(&self) {
let mut metadata = self.lock_metadata();
metadata.swap_primary();
*self.layout.lock().unwrap() = metadata.get_primary_layout();
}
pub(crate) fn begin_repair(&self) -> Result<()> {
let mut metadata = self.lock_metadata();
if !metadata
.primary_slot()
.verify_checksum(metadata.get_checksum_type())
{
metadata.swap_primary();
*self.layout.lock().unwrap() = metadata.get_primary_layout();
assert!(metadata
.primary_slot()
.verify_checksum(metadata.get_checksum_type()));
} else {
let secondary_newer = metadata
.secondary_slot()
.get_last_committed_transaction_id()
> metadata.primary_slot().get_last_committed_transaction_id();
if secondary_newer
&& metadata
.secondary_slot()
.verify_checksum(metadata.get_checksum_type())
{
metadata.swap_primary();
*self.layout.lock().unwrap() = metadata.get_primary_layout();
}
}
let layout = self.layout.lock().unwrap();
metadata.initialize_region_tracker(&layout);
let (mut region_tracker, mut regions) = metadata.allocators_mut(&layout)?;
let num_regions = layout.num_regions();
let mut regional_allocators = vec![];
for i in 0..num_regions {
let mut region = regions.get_region_mut(i);
let region_layout = layout.region_layout(i);
region.initialize(
region_layout.num_pages(),
layout.full_region_layout().num_pages(),
region_layout.max_order(),
);
let allocator = BuddyAllocator::new(
region_layout.num_pages(),
layout.full_region_layout().num_pages(),
region_layout.max_order(),
);
let highest_free = allocator
.highest_free_order(region.allocator_state_mut())
.unwrap();
region_tracker.mark_free(highest_free, i as u64);
regional_allocators.push(allocator);
}
let mut guard = self.regional_allocators.lock().unwrap();
*guard = Some(layout.create_allocators());
Ok(())
}
pub(crate) fn mark_pages_allocated(
&self,
allocated_pages: impl Iterator<Item = PageNumber>,
) -> Result<()> {
let mut metadata = self.lock_metadata();
let layout = self.layout.lock().unwrap();
let (_, mut regions) = metadata.allocators_mut(&layout)?;
let regional_allocators = self.regional_allocators.lock().unwrap();
for page_number in allocated_pages {
let region_index = page_number.region as usize;
let mut region = regions.get_region_mut(region_index);
let mem = region.allocator_state_mut();
regional_allocators.as_ref().unwrap()[region_index].record_alloc(
mem,
page_number.page_index as u64,
page_number.page_order as usize,
);
}
Ok(())
}
pub(crate) fn end_repair(&self) -> Result<()> {
let mut metadata = self.lock_metadata();
self.mmap.flush()?;
metadata.set_recovery(false);
self.mmap.flush()
}
fn lock_metadata(&self) -> MetadataAccessor {
unsafe { MetadataAccessor::new(&self.mmap, self.metadata_guard.lock().unwrap()) }
}
pub(crate) fn commit(
&self,
data_root: Option<(PageNumber, Checksum)>,
freed_root: Option<(PageNumber, Checksum)>,
transaction_id: u64,
eventual: bool,
) -> Result {
#[cfg(debug_assertions)]
debug_assert!(self.open_dirty_pages.lock().unwrap().is_empty());
assert!(self.regional_allocators.lock().unwrap().is_some());
let mut metadata = self.lock_metadata();
let checksum_type = metadata.get_checksum_type();
let mut layout = self.layout.lock().unwrap();
let mut shrunk = false;
if self.dynamic_growth {
shrunk = self.try_shrink(&mut metadata, &mut layout)?;
};
let mut secondary = metadata.secondary_slot_mut();
secondary.set_last_committed_transaction_id(transaction_id);
secondary.set_root_page(data_root);
secondary.set_freed_root(freed_root);
secondary.set_data_section_layout(
layout.num_full_regions(),
layout.trailing_region_layout().map(|x| x.num_pages()),
);
secondary.update_checksum(checksum_type);
if matches!(self.checksum_type, ChecksumType::Zero) {
if eventual {
self.mmap.eventual_flush()?;
} else {
self.mmap.flush()?;
}
}
metadata.swap_primary();
if eventual {
self.mmap.eventual_flush()?;
} else {
self.mmap.flush()?;
}
drop(metadata);
if shrunk {
unsafe {
self.mmap.resize(layout.len().try_into().unwrap())?;
}
}
self.log_since_commit.lock().unwrap().clear();
self.allocated_since_commit.lock().unwrap().clear();
self.read_from_secondary.store(false, Ordering::Release);
Ok(())
}
pub(crate) fn non_durable_commit(
&self,
data_root: Option<(PageNumber, Checksum)>,
freed_root: Option<(PageNumber, Checksum)>,
transaction_id: u64,
) -> Result {
#[cfg(debug_assertions)]
debug_assert!(self.open_dirty_pages.lock().unwrap().is_empty());
assert!(self.regional_allocators.lock().unwrap().is_some());
let mut metadata = self.lock_metadata();
let checksum_type = metadata.get_checksum_type();
let layout = self.layout.lock().unwrap();
let mut secondary = metadata.secondary_slot_mut();
secondary.set_last_committed_transaction_id(transaction_id);
secondary.set_root_page(data_root);
secondary.set_freed_root(freed_root);
secondary.set_data_section_layout(
layout.num_full_regions(),
layout.trailing_region_layout().map(|x| x.num_pages()),
);
secondary.update_checksum(checksum_type);
self.log_since_commit.lock().unwrap().clear();
self.allocated_since_commit.lock().unwrap().clear();
self.read_from_secondary.store(true, Ordering::Release);
Ok(())
}
pub(crate) fn rollback_uncommitted_writes(&self) -> Result {
#[cfg(debug_assertions)]
debug_assert!(self.open_dirty_pages.lock().unwrap().is_empty());
let mut metadata = self.lock_metadata();
let restore = if self.read_from_secondary.load(Ordering::Acquire) {
metadata.get_secondary_layout()
} else {
metadata.get_primary_layout()
};
let mut regional_guard = self.regional_allocators.lock().unwrap();
let mut layout = self.layout.lock().unwrap();
let (mut region_tracker, mut regions) = metadata.allocators_mut(&layout)?;
for op in self.log_since_commit.lock().unwrap().drain(..).rev() {
match op {
AllocationOp::Allocate(page_number) => {
let region_index = page_number.region as usize;
region_tracker.mark_free(page_number.page_order as usize, region_index as u64);
let mut region = regions.get_region_mut(region_index);
let mem = region.allocator_state_mut();
regional_guard.as_ref().unwrap()[region_index].free(
mem,
page_number.page_index as u64,
page_number.page_order as usize,
);
}
AllocationOp::Free(page_number) | AllocationOp::FreeUncommitted(page_number) => {
let region_index = page_number.region as usize;
let mut region = regions.get_region_mut(region_index);
let mem = region.allocator_state_mut();
regional_guard.as_ref().unwrap()[region_index].record_alloc(
mem,
page_number.page_index as u64,
page_number.page_order as usize,
);
}
}
}
self.allocated_since_commit.lock().unwrap().clear();
assert!(restore.len() <= layout.len());
if restore.len() < layout.len() {
regional_guard
.as_mut()
.unwrap()
.drain(restore.num_regions()..);
let last_region_index = restore.num_regions() - 1;
let last_region = restore.region_layout(last_region_index);
let mut region = regions.get_region_mut(last_region_index);
let allocator_data = region.allocator_state_mut();
let last_allocator = &mut regional_guard.as_mut().unwrap()[last_region_index];
last_allocator.resize(allocator_data, last_region.num_pages());
*layout = restore;
unsafe {
self.mmap.resize(layout.len().try_into().unwrap())?;
}
}
Ok(())
}
pub(crate) fn get_page(&self, page_number: PageNumber) -> PageImpl {
#[cfg(debug_assertions)]
{
debug_assert!(
!self.open_dirty_pages.lock().unwrap().contains(&page_number),
"{:?}",
page_number
);
*(self
.read_page_ref_counts
.lock()
.unwrap()
.entry(page_number)
.or_default()) += 1;
}
let mem = unsafe {
self.mmap.get_memory(page_number.address_range(
self.db_header_size,
self.region_size,
self.region_header_with_padding_size,
self.page_size,
))
};
PageImpl {
mem,
page_number,
#[cfg(debug_assertions)]
open_pages: &self.read_page_ref_counts,
}
}
pub(crate) unsafe fn get_page_mut(&self, page_number: PageNumber) -> PageMut {
#[cfg(debug_assertions)]
{
assert!(!self
.read_page_ref_counts
.lock()
.unwrap()
.contains_key(&page_number));
assert!(self.open_dirty_pages.lock().unwrap().insert(page_number));
}
let address_range = page_number.address_range(
self.db_header_size,
self.region_size,
self.region_header_with_padding_size,
self.page_size,
);
let mem = self.mmap.get_memory_mut(address_range);
PageMut {
mem,
page_number,
#[cfg(debug_assertions)]
open_pages: &self.open_dirty_pages,
}
}
pub(crate) fn get_data_root(&self) -> Option<(PageNumber, Checksum)> {
let metadata = self.lock_metadata();
if self.read_from_secondary.load(Ordering::Acquire) {
metadata.secondary_slot().get_root_page()
} else {
metadata.primary_slot().get_root_page()
}
}
pub(crate) fn get_freed_root(&self) -> Option<(PageNumber, Checksum)> {
let metadata = self.lock_metadata();
if self.read_from_secondary.load(Ordering::Acquire) {
metadata.secondary_slot().get_freed_root_page()
} else {
metadata.primary_slot().get_freed_root_page()
}
}
pub(crate) fn get_last_committed_transaction_id(&self) -> Result<u64> {
let metadata = self.lock_metadata();
if self.read_from_secondary.load(Ordering::Acquire) {
Ok(metadata
.secondary_slot()
.get_last_committed_transaction_id())
} else {
Ok(metadata.primary_slot().get_last_committed_transaction_id())
}
}
pub(crate) unsafe fn free(&self, page: PageNumber) -> Result {
let mut mut_page = self.get_page_mut(page);
mut_page.memory_mut().fill(0);
let mut metadata = self.lock_metadata();
let layout = self.layout.lock().unwrap();
let (mut region_tracker, mut regions) = metadata.allocators_mut(&layout)?;
let region_index = page.region as usize;
let mut region = regions.get_region_mut(region_index);
let mem = region.allocator_state_mut();
self.regional_allocators.lock().unwrap().as_ref().unwrap()[region_index].free(
mem,
page.page_index as u64,
page.page_order as usize,
);
region_tracker.mark_free(page.page_order as usize, region_index as u64);
self.log_since_commit
.lock()
.unwrap()
.push(AllocationOp::Free(page));
Ok(())
}
pub(crate) unsafe fn free_if_uncommitted(&self, page: PageNumber) -> Result<bool> {
if self.allocated_since_commit.lock().unwrap().remove(&page) {
let mut mut_page = self.get_page_mut(page);
mut_page.memory_mut().fill(0);
let mut metadata = self.lock_metadata();
let layout = self.layout.lock().unwrap();
let (mut region_tracker, mut regions) = metadata.allocators_mut(&layout)?;
let mut region = regions.get_region_mut(page.region as usize);
let mem = region.allocator_state_mut();
self.regional_allocators.lock().unwrap().as_ref().unwrap()[page.region as usize].free(
mem,
page.page_index as u64,
page.page_order as usize,
);
region_tracker.mark_free(page.page_order as usize, page.region as u64);
self.log_since_commit
.lock()
.unwrap()
.push(AllocationOp::FreeUncommitted(page));
Ok(true)
} else {
Ok(false)
}
}
pub(crate) fn uncommitted(&self, page: PageNumber) -> bool {
self.allocated_since_commit.lock().unwrap().contains(&page)
}
fn allocate_helper(
&self,
metadata: &mut MetadataAccessor,
layout: &DatabaseLayout,
required_order: usize,
) -> Result<PageNumber> {
let regional_guard = self.regional_allocators.lock().unwrap();
let (mut region_tracker, mut regions) = metadata.allocators_mut(layout)?;
loop {
let candidate_region = region_tracker.find_free(required_order)? as usize;
let mut region = regions.get_region_mut(candidate_region);
let mem = region.allocator_state_mut();
match regional_guard.as_ref().unwrap()[candidate_region].alloc(mem, required_order) {
Ok(page) => {
return Ok(PageNumber::new(
candidate_region as u32,
page as u32,
required_order as u8,
));
}
Err(err) => {
if matches!(err, Error::OutOfSpace) {
region_tracker.mark_full(required_order, candidate_region as u64);
} else {
return Err(err);
}
}
}
}
}
fn try_shrink(
&self,
metadata: &mut MetadataAccessor,
layout: &mut DatabaseLayout,
) -> Result<bool> {
let mut allocators = self.regional_allocators.lock().unwrap();
let last_region_index = layout.num_regions() - 1;
let last_allocator = allocators.as_ref().unwrap()[last_region_index].clone();
let last_region = layout.region_layout(last_region_index);
let region = metadata.get_region(last_region_index, layout);
let allocator_data = region.allocator_state();
let trailing_free = last_allocator.trailing_free_pages(allocator_data);
if trailing_free < last_allocator.len() / 2 {
return Ok(false);
}
let reduce_to_pages = if layout.num_regions() > 1 && trailing_free == last_allocator.len() {
0
} else {
max(MIN_USABLE_PAGES, last_allocator.len() - trailing_free)
};
let (mut region_tracker, mut regions) = metadata.allocators_mut(layout)?;
let new_usable_bytes = if reduce_to_pages == 0 {
region_tracker.mark_full(0, last_region_index as u64);
allocators
.as_mut()
.unwrap()
.pop()
.expect("allocators should not be empty");
layout.usable_bytes() - last_region.usable_bytes()
} else {
let mut region = regions.get_region_mut(last_region_index);
let mem = region.allocator_state_mut();
allocators.as_mut().unwrap()[last_region_index].resize(mem, reduce_to_pages);
layout.usable_bytes()
- ((last_allocator.len() - reduce_to_pages) as u64)
* (metadata.get_page_size() as u64)
};
let new_layout = DatabaseLayout::calculate(
metadata.get_max_capacity(),
new_usable_bytes,
metadata.get_region_max_usable_bytes(),
self.page_size,
)?;
assert!(new_layout.len() <= layout.len());
assert_eq!(new_layout.header_bytes(), layout.header_bytes());
assert_eq!(new_layout.header_bytes(), self.db_header_size);
*layout = new_layout;
Ok(true)
}
fn grow(
&self,
metadata: &mut MetadataAccessor,
layout: &mut DatabaseLayout,
required_order_allocation: usize,
) -> Result<()> {
let required_growth =
2u64.pow(required_order_allocation as u32) * metadata.get_page_size() as u64;
let max_region_size = metadata.get_region_max_usable_bytes();
let next_desired_size = if layout.num_full_regions() > 0 {
if let Some(trailing) = layout.trailing_region_layout() {
if 2 * required_growth < max_region_size - trailing.usable_bytes() {
layout.usable_bytes() + (max_region_size - trailing.usable_bytes())
} else {
layout.usable_bytes() + max_region_size
}
} else {
layout.usable_bytes() + max_region_size
}
} else {
max(
layout.usable_bytes() * 2,
layout.usable_bytes() + required_growth * 2,
)
};
let new_layout = DatabaseLayout::calculate(
metadata.get_max_capacity(),
next_desired_size,
metadata.get_region_max_usable_bytes(),
self.page_size,
)?;
assert!(new_layout.len() >= layout.len());
assert_eq!(new_layout.header_bytes(), layout.header_bytes());
assert_eq!(new_layout.header_bytes(), self.db_header_size);
if new_layout.len() == layout.len() {
return Err(Error::OutOfSpace);
}
unsafe {
self.mmap.resize(new_layout.len().try_into().unwrap())?;
}
let mut allocators = self.regional_allocators.lock().unwrap();
let mut new_allocators = vec![];
for i in 0..new_layout.num_regions() {
let new_region_base = new_layout.region_base_address(i);
let new_region = new_layout.region_layout(i);
let new_allocator = if i < layout.num_regions() {
let old_region_base = layout.region_base_address(i);
let old_region = layout.region_layout(i);
assert_eq!(old_region_base, new_region_base);
let mut allocator = allocators.as_ref().unwrap()[i].clone();
if new_region.len() != old_region.len() {
let (mut region_tracker, mut regions) = metadata.allocators_mut(&new_layout)?;
let mut region = regions.get_region_mut(i);
let mem = region.allocator_state_mut();
allocator.resize(mem, new_region.num_pages());
let highest_free = allocator.highest_free_order(mem).unwrap();
region_tracker.mark_free(highest_free, i as u64);
}
allocator
} else {
let (mut region_tracker, mut regions) = metadata.allocators_mut(&new_layout)?;
let mut region = regions.get_region_mut(i);
region.initialize(
new_region.num_pages(),
new_layout.full_region_layout().num_pages(),
new_region.max_order(),
);
let allocator = BuddyAllocator::new(
new_region.num_pages(),
new_layout.full_region_layout().num_pages(),
new_region.max_order(),
);
let mem = region.allocator_state_mut();
let highest_free = allocator.highest_free_order(mem).unwrap();
region_tracker.mark_free(highest_free, i as u64);
allocator
};
new_allocators.push(new_allocator);
}
*allocators = Some(new_allocators);
*layout = new_layout;
Ok(())
}
pub(crate) fn allocate(&self, allocation_size: usize) -> Result<PageMut> {
let required_pages = (allocation_size + self.page_size - 1) / self.page_size;
let required_order = ceil_log2(required_pages);
let mut metadata = self.lock_metadata();
let max_capacity = metadata.get_max_capacity();
let mut layout = self.layout.lock().unwrap();
let page_number = match self.allocate_helper(&mut metadata, &layout, required_order) {
Ok(page_number) => page_number,
Err(err) => {
if matches!(err, Error::OutOfSpace) && (layout.len() as u64) < max_capacity {
self.grow(&mut metadata, &mut layout, required_order)?;
self.allocate_helper(&mut metadata, &layout, required_order)?
} else {
return Err(err);
}
}
};
self.allocated_since_commit
.lock()
.unwrap()
.insert(page_number);
self.log_since_commit
.lock()
.unwrap()
.push(AllocationOp::Allocate(page_number));
#[cfg(debug_assertions)]
{
assert!(!self
.read_page_ref_counts
.lock()
.unwrap()
.contains_key(&page_number));
assert!(self.open_dirty_pages.lock().unwrap().insert(page_number));
}
let address_range = page_number.address_range(
self.db_header_size,
self.region_size,
self.region_header_with_padding_size,
self.page_size,
);
let mem = unsafe { self.mmap.get_memory_mut(address_range) };
debug_assert!(mem.len() >= allocation_size);
Ok(PageMut {
mem,
page_number,
#[cfg(debug_assertions)]
open_pages: &self.open_dirty_pages,
})
}
pub(crate) fn count_free_pages(&self) -> Result<usize> {
let mut metadata = self.lock_metadata();
let regional_guard = self.regional_allocators.lock().unwrap();
let layout = self.layout.lock().unwrap();
let mut count = 0;
for i in 0..layout.num_regions() {
let region = metadata.get_region(i, &layout);
let mem = region.allocator_state();
count += regional_guard.as_ref().unwrap()[i].count_free_pages(mem);
}
let max_layout = DatabaseLayout::calculate(
metadata.get_max_capacity(),
metadata.get_max_capacity(),
metadata.get_region_max_usable_bytes(),
self.page_size,
)
.unwrap();
let potential_growth_pages: usize = ((max_layout.usable_bytes() - layout.usable_bytes())
/ (self.page_size as u64))
.try_into()
.unwrap();
Ok(count + potential_growth_pages)
}
pub(crate) fn get_page_size(&self) -> usize {
self.page_size
}
}
impl Drop for TransactionalMemory {
fn drop(&mut self) {
if self.read_from_secondary.load(Ordering::Acquire) {
if let Ok(non_durable_transaction_id) = self.get_last_committed_transaction_id() {
let root = self.get_data_root();
let freed_root = self.get_freed_root();
if self
.commit(root, freed_root, non_durable_transaction_id, false)
.is_err()
{
eprintln!(
"Failure while finalizing non-durable commit. Database may have rolled back"
);
}
} else {
eprintln!(
"Failure while finalizing non-durable commit. Database may have rolled back"
);
}
}
match self.regional_allocators.lock() {
Ok(allocators) => {
if self.mmap.flush().is_ok() && allocators.is_some() {
self.lock_metadata().set_recovery(false);
let _ = self.mmap.flush();
}
}
Err(_) => {
let _ = self.mmap.flush();
eprintln!("Failure while closing database");
}
}
}
}
#[cfg(test)]
mod test {
use crate::db::TableDefinition;
use crate::tree_store::page_store::page_manager::{
DB_HEADER_SIZE, GOD_BYTE_OFFSET, MAGICNUMBER, MIN_USABLE_PAGES, PRIMARY_BIT,
RECOVERY_REQUIRED, ROOT_CHECKSUM_OFFSET, TRANSACTION_0_OFFSET, TRANSACTION_1_OFFSET,
};
use crate::tree_store::page_store::utils::get_page_size;
use crate::tree_store::page_store::TransactionalMemory;
use crate::{Database, Error, ReadableTable, WriteStrategy};
use std::fs::OpenOptions;
use std::io::{Read, Seek, SeekFrom, Write};
use std::mem::size_of;
use tempfile::NamedTempFile;
const X: TableDefinition<[u8], [u8]> = TableDefinition::new("x");
#[test]
fn repair_allocator_no_checksums() {
let tmpfile: NamedTempFile = NamedTempFile::new().unwrap();
let max_size = 1024 * 1024;
let db = unsafe {
Database::builder()
.set_write_strategy(WriteStrategy::TwoPhase)
.create(tmpfile.path(), max_size)
.unwrap()
};
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(X).unwrap();
table.insert(b"hello", b"world").unwrap();
}
write_txn.commit().unwrap();
let write_txn = db.begin_write().unwrap();
let free_pages = write_txn.stats().unwrap().free_pages();
write_txn.abort().unwrap();
drop(db);
let mut file = OpenOptions::new()
.read(true)
.write(true)
.open(tmpfile.path())
.unwrap();
file.seek(SeekFrom::Start(GOD_BYTE_OFFSET as u64)).unwrap();
let mut buffer = [0u8; 1];
file.read_exact(&mut buffer).unwrap();
file.seek(SeekFrom::Start(GOD_BYTE_OFFSET as u64)).unwrap();
buffer[0] |= RECOVERY_REQUIRED;
file.write_all(&buffer).unwrap();
assert!(TransactionalMemory::new(
file,
max_size as u64,
None,
None,
true,
Some(WriteStrategy::TwoPhase)
)
.unwrap()
.needs_repair()
.unwrap());
let db2 = unsafe {
Database::builder()
.set_write_strategy(WriteStrategy::TwoPhase)
.create(tmpfile.path(), max_size)
.unwrap()
};
let write_txn = db2.begin_write().unwrap();
assert_eq!(free_pages, write_txn.stats().unwrap().free_pages());
{
let mut table = write_txn.open_table(X).unwrap();
table.insert(b"hello2", b"world2").unwrap();
}
write_txn.commit().unwrap();
}
#[test]
fn repair_allocator_checksums() {
let tmpfile: NamedTempFile = NamedTempFile::new().unwrap();
let max_size = 1024 * 1024;
let db = unsafe {
Database::builder()
.set_write_strategy(WriteStrategy::Checksum)
.create(tmpfile.path(), max_size)
.unwrap()
};
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(X).unwrap();
table.insert(b"hello", b"world").unwrap();
}
write_txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(X).unwrap();
table.insert(b"hello", b"world2").unwrap();
}
write_txn.commit().unwrap();
drop(read_txn);
drop(db);
let mut file = OpenOptions::new()
.read(true)
.write(true)
.open(tmpfile.path())
.unwrap();
file.seek(SeekFrom::Start(GOD_BYTE_OFFSET as u64)).unwrap();
let mut buffer = [0u8; 1];
file.read_exact(&mut buffer).unwrap();
file.seek(SeekFrom::Start(GOD_BYTE_OFFSET as u64)).unwrap();
buffer[0] |= RECOVERY_REQUIRED;
file.write_all(&buffer).unwrap();
let primary_slot_offset = if buffer[0] & PRIMARY_BIT == 0 {
TRANSACTION_0_OFFSET
} else {
TRANSACTION_1_OFFSET
};
file.seek(SeekFrom::Start(
(primary_slot_offset + ROOT_CHECKSUM_OFFSET) as u64,
))
.unwrap();
file.write_all(&[0; size_of::<u128>()]).unwrap();
assert!(TransactionalMemory::new(
file,
max_size as u64,
None,
None,
true,
Some(WriteStrategy::Checksum)
)
.unwrap()
.needs_repair()
.unwrap());
let db2 = unsafe { Database::create(tmpfile.path(), max_size).unwrap() };
let write_txn = db2.begin_write().unwrap();
{
let mut table = write_txn.open_table(X).unwrap();
assert_eq!(table.get(b"hello").unwrap().unwrap(), b"world");
table.insert(b"hello2", b"world2").unwrap();
}
write_txn.commit().unwrap();
}
#[test]
fn repair_insert_reserve_regression() {
let tmpfile: NamedTempFile = NamedTempFile::new().unwrap();
let max_size = 1024 * 1024;
let db = unsafe {
Database::builder()
.set_write_strategy(WriteStrategy::Checksum)
.create(tmpfile.path(), max_size)
.unwrap()
};
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(X).unwrap();
let mut value = table.insert_reserve(b"hello", 5).unwrap();
value.as_mut().copy_from_slice(b"world");
}
write_txn.commit().unwrap();
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(X).unwrap();
let mut value = table.insert_reserve(b"hello2", 5).unwrap();
value.as_mut().copy_from_slice(b"world");
}
write_txn.commit().unwrap();
drop(db);
let mut file = OpenOptions::new()
.read(true)
.write(true)
.open(tmpfile.path())
.unwrap();
file.seek(SeekFrom::Start(GOD_BYTE_OFFSET as u64)).unwrap();
let mut buffer = [0u8; 1];
file.read_exact(&mut buffer).unwrap();
file.seek(SeekFrom::Start(GOD_BYTE_OFFSET as u64)).unwrap();
buffer[0] |= RECOVERY_REQUIRED;
file.write_all(&buffer).unwrap();
assert!(TransactionalMemory::new(
file,
max_size as u64,
None,
None,
true,
Some(WriteStrategy::Checksum)
)
.unwrap()
.needs_repair()
.unwrap());
unsafe { Database::open(tmpfile.path()).unwrap() };
}
#[test]
fn too_small_db() {
let tmpfile: NamedTempFile = NamedTempFile::new().unwrap();
let result = unsafe { Database::create(tmpfile.path(), 1) };
assert!(matches!(result, Err(Error::OutOfSpace)));
let tmpfile: NamedTempFile = NamedTempFile::new().unwrap();
let result = unsafe { Database::create(tmpfile.path(), 1024) };
assert!(matches!(result, Err(Error::OutOfSpace)));
let tmpfile: NamedTempFile = NamedTempFile::new().unwrap();
let result =
unsafe { Database::create(tmpfile.path(), MIN_USABLE_PAGES * get_page_size() - 1) };
assert!(matches!(result, Err(Error::OutOfSpace)));
}
#[test]
fn smallest_db() {
let tmpfile: NamedTempFile = NamedTempFile::new().unwrap();
unsafe {
Database::create(
tmpfile.path(),
DB_HEADER_SIZE + (MIN_USABLE_PAGES + 2) * get_page_size(),
)
.unwrap();
}
}
#[test]
fn magic_number() {
assert!(std::str::from_utf8(&MAGICNUMBER).is_err());
assert!(MAGICNUMBER.iter().any(|x| *x & 0x80 != 0));
assert!(MAGICNUMBER.iter().any(|x| *x < 0x20 || *x > 0x7E));
assert!(MAGICNUMBER.iter().any(|x| *x >= 0x20 && *x <= 0x7E));
assert!(MAGICNUMBER.iter().any(|x| *x >= 0xA0));
assert!(MAGICNUMBER.iter().any(|x| *x < 0x09
|| *x == 0x0B
|| (0x0E <= *x && *x <= 0x1F)
|| (0x7F <= *x && *x <= 0x9F)));
}
}