use std::{
ops::Deref,
slice::from_raw_parts_mut,
sync::{Arc, atomic::AtomicU64},
};
use log::{debug, info};
use parking_lot::{Mutex, RwLockReadGuard};
use wdev::Device;
use wepoch::LightEpoch;
use wram::AlignedBuf;
use wrecord::{
ADDRESS_MASK, HEADER_SIZE, RecordHeader, RecordRef, checked_record_size, encode_to_slice,
};
use crate::{
address::{AddressManager, AddressSnapshot},
buffer::CircularPageBuffer,
config::HybridLogConfig,
error::{Error, Result},
flush::PendingFlushList,
};
mod append;
mod inplace;
mod io;
mod shift;
pub const PAD_KEY_LEN: u32 = u32::MAX;
const DISK_READ_CACHE_SLOTS: usize = 2;
const DISK_READ_CACHE_MASK: usize = DISK_READ_CACHE_SLOTS - 1;
const DISK_READ_PROBE_LEN: usize = 4096;
type DiskPageSlot = Option<(u64, AlignedBuf)>;
pub struct HybridLog<D: Device> {
pub config: HybridLogConfig,
pub device: Arc<D>,
pub epoch: Arc<LightEpoch>,
pub addresses: Arc<AddressManager>,
pub buffer: CircularPageBuffer,
page_turn_lock: Mutex<()>,
flush_staging: Mutex<Option<AlignedBuf>>,
disk_read_cache: Mutex<Box<[DiskPageSlot]>>,
last_probe_page: AtomicU64,
pub pending_flush: PendingFlushList,
}
pub(crate) enum PageBytes<'a> {
Disk(AlignedBuf),
Raw(&'a [u8]),
Locked(RwLockReadGuard<'a, AlignedBuf>),
}
impl Deref for PageBytes<'_> {
type Target = [u8];
#[inline]
fn deref(&self) -> &[u8] {
match self {
Self::Disk(buf) => buf.as_slice(),
Self::Raw(slice) => slice,
Self::Locked(guard) => guard.as_slice(),
}
}
}
#[inline]
pub(crate) const fn reject_pad(header: RecordHeader, addr: u64) -> Result<()> {
if header.key_len == PAD_KEY_LEN
|| (header.key_len == 0 && header.val_len == 0 && header.prev_address == 0)
{
Err(Error::PadRecord(addr))
} else {
Ok(())
}
}
#[inline]
pub(crate) fn parse_record_from_slice(
page_slice: &[u8],
offset: usize,
addr: u64,
page_size: usize,
) -> Result<RecordRef<'_>> {
if offset + HEADER_SIZE > page_slice.len() || offset + HEADER_SIZE > page_size {
return Err(Error::PadRecord(addr));
}
let header = RecordHeader::from_slice(&page_slice[offset..offset + HEADER_SIZE])?;
reject_pad(header, addr)?;
let logical_size = header.record_size();
let physical_size = header.physical_size();
let end = match offset.checked_add(logical_size) {
Some(e)
if offset.saturating_add(physical_size) <= page_slice.len()
&& offset.saturating_add(physical_size) <= page_size =>
{
e
}
_ => {
return Err(Error::RecordCorrupted {
addr,
detail: "记录完整内容超出页面容量边界".into(),
});
}
};
let key_start = offset + HEADER_SIZE;
let key_end = key_start + header.key_len as usize;
let (key, value) = unsafe {
(
page_slice.get_unchecked(key_start..key_end),
page_slice.get_unchecked(key_end..end),
)
};
Ok(RecordRef { header, key, value })
}
pub(crate) struct RecParams<'a> {
pub prev_addr: u64,
pub key: &'a [u8],
pub val: &'a [u8],
pub is_tombstone: bool,
}
impl<D: Device> HybridLog<D> {
fn assemble(
config: HybridLogConfig,
device: Arc<D>,
epoch: Arc<LightEpoch>,
addresses: Arc<AddressManager>,
buffer: CircularPageBuffer,
) -> Self {
Self {
config,
device,
epoch,
addresses,
buffer,
page_turn_lock: Mutex::new(()),
flush_staging: Mutex::new(None),
disk_read_cache: Mutex::new(vec![None; DISK_READ_CACHE_SLOTS].into_boxed_slice()),
last_probe_page: AtomicU64::new(0),
pending_flush: PendingFlushList::new(),
}
}
pub fn new(config: HybridLogConfig, device: Arc<D>, epoch: Arc<LightEpoch>) -> Result<Self> {
let initial_address = config.initial_address;
let addresses = Arc::new(AddressManager::new(initial_address));
let buffer = CircularPageBuffer::new(&config)?;
let initial_page = config.page_id(initial_address);
buffer.clear_page(initial_page);
info!(
"初始化 HybridLog: page_size={}, num_pages={}, mutable_fraction={}, initial_address={initial_address:#x}",
config.page_size, config.num_pages, config.mutable_fraction
);
Ok(Self::assemble(config, device, epoch, addresses, buffer))
}
pub async fn recover(
config: HybridLogConfig,
device: Arc<D>,
epoch: Arc<LightEpoch>,
snapshot: AddressSnapshot,
) -> Result<Self> {
if !snapshot.validate() {
let mut ibuf = itoa::Buffer::new();
let mut msg = String::from("恢复快照违反单调不变式: ");
for (name, v) in [
("begin", snapshot.begin),
("safe_head", snapshot.safe_head),
("head", snapshot.head),
("safe_ro", snapshot.safe_read_only),
("ro", snapshot.read_only),
("tail", snapshot.tail),
("flushed_until", snapshot.flushed_until),
] {
msg.push_str(name);
msg.push('=');
msg.push_str(ibuf.format(v));
msg.push_str(", ");
}
msg.pop();
msg.pop();
return Err(Error::InvalidState(msg));
}
let page_size = config.page_size;
let tail_page = config.page_id(snapshot.tail);
if snapshot.tail > snapshot.head
&& tail_page >= config.page_id(snapshot.head) + config.num_pages as u64
{
let mut ibuf = itoa::Buffer::new();
let mut msg = String::from("恢复快照驻留窗口超过环形页数: head_page=");
msg.push_str(ibuf.format(config.page_id(snapshot.head)));
msg.push_str(", tail_page=");
msg.push_str(ibuf.format(tail_page));
msg.push_str(", num_pages=");
msg.push_str(ibuf.format(config.num_pages));
return Err(Error::InvalidState(msg));
}
let buffer = CircularPageBuffer::new(&config)?;
let addresses = Arc::new(AddressManager::with_snapshot(snapshot));
if snapshot.tail > snapshot.head {
let head_page = config.page_id(snapshot.head);
let first_start = config.page_start_address(head_page);
let span = (config.page_start_address(tail_page.saturating_add(1)) - first_start) as usize;
let span_buf = match device.read_range(first_start, span).await {
Ok(buf) => Some(buf),
Err(e) => {
debug!("恢复阶段整段批量预热失败,退回逐页加载: {e}");
None
}
};
for p in head_page..=tail_page {
let page_start = config.page_start_address(p);
if page_start >= snapshot.flushed_until {
buffer.clear_page(p);
continue;
}
match span_buf.as_ref() {
Some(all) => {
let base = (page_start - first_start) as usize;
let end = (base + page_size).min(all.len());
if base < all.len() {
buffer.load_page(p, &all[base..end]);
} else {
buffer.clear_page(p);
}
}
None => match device.read_range(page_start, page_size).await {
Ok(buf) => buffer.load_page(p, &buf),
Err(e) => {
debug!("恢复阶段预热加载逻辑页 {p} 失败(空设备或短文件,按空页处理): {e}");
buffer.clear_page(p);
continue;
}
},
}
let scrub = snapshot
.flushed_until
.saturating_sub(page_start)
.min(page_size as u64) as usize;
if scrub < page_size {
buffer.clear_page_from_offset(p, scrub);
}
}
} else {
buffer.clear_page(tail_page);
}
info!(
"从快照恢复 HybridLog: page_size={}, tail={:#x}, head={:#x}, flushed_until={:#x}",
config.page_size, snapshot.tail, snapshot.head, snapshot.flushed_until
);
Ok(Self::assemble(config, device, epoch, addresses, buffer))
}
pub(crate) fn probe_resident(&self, addr: u64) -> Result<Option<PageBytes<'_>>> {
let head = self.addresses.head();
if addr < head || addr >= self.addresses.tail() {
return Ok(None);
}
let page_id = self.config.page_id(addr);
if let Some(page_slice) = unsafe { self.buffer.try_read_page_unlocked(page_id) }
&& addr >= self.addresses.head()
{
return Ok(Some(PageBytes::Raw(page_slice)));
}
let guard = self.buffer.read_page(page_id);
if addr >= self.addresses.head() && self.buffer.is_page_loaded(page_id) {
return Ok(Some(PageBytes::Locked(guard)));
}
Ok(None)
}
#[inline(always)]
pub unsafe fn get_physical_address(&self, addr: u64) -> *const u8 {
unsafe { self.buffer.get_physical_address(addr) }
}
#[inline]
pub(crate) fn validate_append_args(&self, p: &RecParams<'_>) -> Result<usize> {
if p.prev_addr & !ADDRESS_MASK != 0 {
return Err(Error::InvalidAddress(p.prev_addr));
}
let rec_size = match checked_record_size(p.key.len(), p.val.len()) {
Some(s) => s,
None => {
return Err(Error::RecordTooLarge {
size: usize::MAX,
page_size: self.config.page_size,
});
}
};
if rec_size > self.config.page_size {
return Err(Error::RecordTooLarge {
size: rec_size,
page_size: self.config.page_size,
});
}
Ok(rec_size)
}
#[inline]
pub(crate) unsafe fn encode_at(
&self,
page_id: u64,
offset: usize,
rec_size: usize,
p: &RecParams<'_>,
) -> Result<()> {
unsafe {
let slot = self.buffer.page_idx(page_id);
let page_ptr = self.buffer.raw_page_ptr_mut(slot);
let dest = from_raw_parts_mut(page_ptr.add(offset), rec_size);
encode_to_slice(dest, p.prev_addr, p.key, p.val, p.is_tombstone)?;
}
Ok(())
}
pub(crate) unsafe fn write_pad_tail(&self, page_id: u64, offset: usize, remaining: usize) {
unsafe {
let slot = self.buffer.page_idx(page_id);
let page_ptr = self.buffer.raw_page_ptr_mut(slot);
let dest = from_raw_parts_mut(page_ptr.add(offset), remaining);
if remaining >= HEADER_SIZE {
let pad_header = RecordHeader {
prev_address: 0,
key_len: PAD_KEY_LEN,
val_len: (remaining - HEADER_SIZE) as u32,
};
dest[..HEADER_SIZE].copy_from_slice(&pad_header.to_bytes());
} else {
dest.fill(0xFF);
}
}
}
}