use std::{
fs::OpenOptions,
io::{Seek, SeekFrom, Write},
path::{Path, PathBuf},
sync::{
Arc,
atomic::{AtomicBool, AtomicUsize, Ordering},
},
thread,
time::{Duration, Instant},
};
use aok::Result;
use wcpr::{CprRecover, CprStore, RecoveredCheckpoint, StoreMeta};
use wdev::SegmentedDevice;
use wepoch::{LightEpoch, Participant};
use whlog::{HybridLog, HybridLogConfig};
use windex::{HashBucket, HashBucketEntry, HashIndex};
pub(crate) const INDEX_BUCKETS: usize = 64;
pub(crate) const PAGE_SIZE: usize = 16 * 1024;
pub(crate) const NUM_PAGES: usize = 8;
pub(crate) const MUTABLE_FRACTION: f64 = 0.5;
pub(crate) const MAX_SESSIONS: usize = 64;
pub(crate) const RC_VIRTUAL_BASE: u64 = 1 << 40;
pub(crate) struct MiniStore {
pub device: Arc<SegmentedDevice>,
pub hlog: Arc<HybridLog<SegmentedDevice>>,
pub index: Arc<HashIndex>,
pub epoch: Arc<LightEpoch>,
pub db_path: PathBuf,
pub meta: Option<wcpr::CheckpointMeta>,
entries: AtomicUsize,
}
impl MiniStore {
pub fn open(db_path: impl AsRef<Path>) -> Result<Arc<Self>> {
let db_path = db_path.as_ref().to_path_buf();
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let config = HybridLogConfig::new(PAGE_SIZE, NUM_PAGES, MUTABLE_FRACTION)?;
let epoch = Arc::new(LightEpoch::new(MAX_SESSIONS));
let hlog = Arc::new(HybridLog::new(
config,
Arc::clone(&device),
Arc::clone(&epoch),
)?);
let index = Arc::new(HashIndex::new(INDEX_BUCKETS)?);
Ok(Arc::new(Self {
device,
hlog,
index,
epoch,
db_path,
meta: None,
entries: AtomicUsize::new(0),
}))
}
pub fn session(&self) -> Result<Participant> {
Ok(self.epoch.register()?)
}
pub fn entry_count(&self) -> usize {
self.entries.load(Ordering::Acquire)
}
async fn append(
&self,
p: &Participant,
key: &[u8],
val: &[u8],
prev: u64,
tombstone: bool,
) -> Result<u64> {
loop {
let attempt = {
let _guard = p.enter();
self.hlog.append(key, val, prev, tombstone)
};
match attempt {
Ok(addr) => return Ok(addr),
Err(whlog::Error::PageNotReady(page_id)) => self.evict_page(page_id).await?,
Err(e) => return Err(e.into()),
}
}
}
async fn evict_page(&self, page_id: u64) -> Result<()> {
let num_pages = self.hlog.config.num_pages as u64;
if page_id < num_pages {
return Ok(());
}
let old_page = page_id - num_pages;
self.hlog.flush_page(old_page).await?;
self.device.sync().await?;
let min_evicted = self.hlog.config.page_start_address(old_page + 1);
self.hlog.shift_read_only_address(min_evicted);
self.hlog.shift_head_address(min_evicted);
self.epoch.bump_epoch();
Ok(())
}
pub fn resolve_main(&self, slot: u64) -> u64 {
let abs = slot & HashBucketEntry::ADDRESS_MASK;
if abs & HashBucketEntry::READ_CACHE_BIT != 0 {
(abs & !HashBucketEntry::READ_CACHE_BIT) - RC_VIRTUAL_BASE
} else {
abs
}
}
pub async fn put(&self, p: &Participant, key: &[u8], val: &[u8]) -> Result<u64> {
let old = {
let _guard = p.enter();
self.index.find_tag(key)
};
let prev = old.map_or(0, |slot| self.resolve_main(slot));
let addr = self.append(p, key, val, prev, false).await?;
{
let _guard = p.enter();
match old {
Some(slot) => {
if !self.index.update_address(key, slot, addr) {
aok::bail!("fixture put: 索引 CAS 更新失败(测试串行场景不应发生): key={key:?}");
}
}
None => self.index.insert(key, addr)?,
}
}
self.entries.fetch_add(1, Ordering::AcqRel);
Ok(addr)
}
pub async fn del(&self, p: &Participant, key: &[u8]) -> Result<u64> {
let old = {
let _guard = p.enter();
self.index.find_tag(key)
};
if old.is_none() {
aok::bail!("fixture del: 键不存在: {key:?}");
}
let prev = old.map_or(0, |slot| self.resolve_main(slot));
let addr = self.append(p, key, &[], prev, true).await?;
{
let _guard = p.enter();
let slot = old.expect("old 已判非空");
if !self.index.update_address(key, slot, addr) {
aok::bail!("fixture del: 索引 CAS 更新失败: key={key:?}");
}
}
self.entries.fetch_sub(1, Ordering::AcqRel);
Ok(addr)
}
pub async fn get(&self, p: &Participant, key: &[u8]) -> Result<Option<Vec<u8>>> {
let slot = {
let _guard = p.enter();
self.index.find_tag(key)
};
let Some(slot) = slot else {
return Ok(None);
};
let addr = self.resolve_main(slot);
if addr == 0 {
return Ok(None);
}
let rec = self.hlog.read_record(addr).await?;
if !rec.key().is_ok_and(|k| k == key) || rec.is_tombstone().unwrap_or(true) {
return Ok(None);
}
Ok(Some(rec.value()?.to_vec()))
}
pub fn install_read_cache_entry(&self, p: &Participant, key: &[u8]) -> Result<()> {
let _guard = p.enter();
let main = self
.index
.find_tag(key)
.ok_or_else(|| aok::anyhow!("install_read_cache_entry: 键不存在: {key:?}"))?;
let bucket = self.index.bucket_for_key(key);
for slot in &bucket.entries[..HashBucket::DATA_ENTRIES] {
let cur = slot.load(Ordering::Acquire);
if cur & HashBucketEntry::ADDRESS_MASK == main {
let virt = (main + RC_VIRTUAL_BASE) | HashBucketEntry::READ_CACHE_BIT;
let new_raw = (cur & !HashBucketEntry::ADDRESS_MASK) | virt;
if slot
.compare_exchange(cur, new_raw, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
{
return Ok(());
}
}
}
aok::bail!("install_read_cache_entry: 未定位到匹配槽位: key={key:?}");
}
pub async fn scan_count(&self, p: &Participant, from: u64, until: u64) -> Result<usize> {
let _guard = p.enter();
let mut iter = self.hlog.scan_iter(from, until);
let mut buf = Vec::with_capacity(PAGE_SIZE);
let mut n = 0usize;
while iter.next_into(&mut buf).await?.is_some() {
n += 1;
}
Ok(n)
}
pub fn scorch_device_beyond_tail(&self, len: usize) -> Result<u64> {
let tail = self.hlog.tail_address();
let mut f = OpenOptions::new().write(true).open(&self.db_path)?;
f.seek(SeekFrom::Start(tail))?;
f.write_all(&vec![0xA5u8; len])?;
f.sync_all()?;
Ok(tail)
}
}
impl CprStore for MiniStore {
type Device = SegmentedDevice;
fn hlog(&self) -> &HybridLog<Self::Device> {
&self.hlog
}
fn index(&self) -> &HashIndex {
&self.index
}
fn epoch(&self) -> &LightEpoch {
&self.epoch
}
fn tail_address(&self) -> u64 {
self.hlog.tail_address()
}
fn begin_address(&self) -> u64 {
self.hlog.begin_address()
}
fn head_address(&self) -> u64 {
self.hlog.head_address()
}
fn shift_read_only_address(&self, target: u64) {
self.hlog.shift_read_only_address(target);
}
async fn flush_all(&self) -> wcpr::Result<()> {
self.hlog.flush_all().await?;
self.device.sync().await?;
Ok(())
}
fn entry_count(&self) -> usize {
self.entries.load(Ordering::Acquire)
}
fn skip_read_cache(&self, addr: u64) -> u64 {
self.resolve_main(addr)
}
fn take_range_index_checkpoints(&self, _dir: &Path, _token: u128) -> wcpr::Result<usize> {
Ok(0)
}
fn take_bftree_checkpoint(&self, _dir: &Path, _token: u128) -> wcpr::Result<usize> {
Ok(0)
}
fn checkpoint_store_meta(&self) -> StoreMeta {
StoreMeta {
index_size: INDEX_BUCKETS,
page_size: PAGE_SIZE,
num_pages: NUM_PAGES,
mutable_fraction: MUTABLE_FRACTION,
max_sessions: MAX_SESSIONS,
enable_revivification: false,
enable_read_cache: true,
read_cache_num_pages: 8,
range_index_dir: None,
bftree_path: None,
next_key_id: 0,
}
}
}
impl CprRecover for MiniStore {
async fn from_recovered(
recovered: RecoveredCheckpoint<Self::Device>,
_checkpoint_dir: &Path,
device: Arc<Self::Device>,
) -> wcpr::Result<Self> {
let db_path = device.segment_path(device.start_segment());
let entry_count = recovered.meta.index_meta.entry_count;
Ok(Self {
db_path,
device: Arc::clone(&device),
meta: Some(recovered.meta),
hlog: recovered.hlog,
index: recovered.index,
epoch: recovered.epoch,
entries: AtomicUsize::new(entry_count),
})
}
}
pub(crate) struct Watchdog {
done: Arc<AtomicBool>,
}
impl Watchdog {
pub fn start(name: &'static str, budget: Duration) -> Self {
let done = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&done);
let start = Instant::now();
let spawned = thread::Builder::new()
.name(format!("watchdog-{name}"))
.spawn(move || {
let mut last_report = 0u64;
while !flag.load(Ordering::Relaxed) {
thread::sleep(Duration::from_secs(1));
let elapsed = start.elapsed().as_secs();
if elapsed >= budget.as_secs() && elapsed > last_report {
last_report = elapsed;
eprintln!(
"[watchdog] 测试 {name} 已运行 {elapsed}s 超预算 {}s 仍未结束,疑似挂死",
budget.as_secs()
);
}
}
});
if spawned.is_err() {
eprintln!("[watchdog] 看门狗线程创建失败,仅损失诊断能力: {name}");
}
Self { done }
}
}
impl Drop for Watchdog {
fn drop(&mut self) {
self.done.store(true, Ordering::Relaxed);
}
}