use pagedb::vfs::memory::MemVfs;
use pagedb::{Db, OpenOptions, PagedbError, ReaderStallPolicy, RealmId};
const KEK: [u8; 32] = [7u8; 32];
const REALM: RealmId = RealmId::new([1u8; 16]);
const LOW_THRESHOLD: u64 = 5;
fn opts_low_threshold() -> OpenOptions {
let mut opts = OpenOptions::default().with_buffer_pool_pages(256);
opts.reader_stall_threshold_pages = LOW_THRESHOLD;
opts
}
async fn fill_until_stall(db: &Db<MemVfs>) -> Option<PagedbError> {
for round in 0..30u32 {
let prefix = format!("k{round:04}");
let mut w = db.begin_write().await.unwrap();
for i in 0u64..8 {
let key = format!("{prefix}_{i:04}");
w.put(key.as_bytes(), &[0u8; 200]).await.unwrap();
}
if let Err(e) = w.commit().await {
return Some(e);
}
let mut w2 = db.begin_write().await.unwrap();
for i in 0u64..8 {
let key = format!("{prefix}_{i:04}");
w2.put(key.as_bytes(), b"x").await.unwrap();
}
if let Err(e) = w2.commit().await {
return Some(e);
}
}
None
}
#[tokio::test(flavor = "current_thread")]
async fn unbounded_never_aborts() {
let vfs = MemVfs::new();
let opts = opts_low_threshold();
let db = Db::open(vfs, KEK, 4096, REALM, opts).await.unwrap();
db.set_reader_stall_policy(ReaderStallPolicy::Unbounded);
let reader = db.begin_read().await.unwrap();
let err = fill_until_stall(&db).await;
assert!(
err.is_none(),
"Unbounded should not produce errors; got {err:?}"
);
let result = reader.get(b"anything").await;
assert!(result.is_ok(), "reader should still be alive: {result:?}");
drop(reader);
}
#[tokio::test(flavor = "current_thread")]
async fn reject_returns_backlog_error() {
let vfs = MemVfs::new();
let opts = opts_low_threshold();
let db = Db::open(vfs, KEK, 4096, REALM, opts).await.unwrap();
db.set_reader_stall_policy(ReaderStallPolicy::Reject);
let _reader = db.begin_read_non_abortable().await.unwrap();
let err = fill_until_stall(&db).await;
assert!(
matches!(err, Some(PagedbError::DeferredFreeBacklog { .. })),
"expected DeferredFreeBacklog, got: {err:?}"
);
}
#[tokio::test(flavor = "current_thread")]
async fn abort_oldest_aborts_old_reader() {
let vfs = MemVfs::new();
let opts = opts_low_threshold();
let db = Db::open(vfs, KEK, 4096, REALM, opts).await.unwrap();
db.set_reader_stall_policy(ReaderStallPolicy::AbortOldest);
{
let mut txn = db.begin_write().await.unwrap();
txn.put(b"hello", b"world").await.unwrap();
txn.commit().await.unwrap();
}
let r1 = db.begin_read().await.unwrap();
let r2 = db.begin_read().await.unwrap();
fill_until_stall(&db).await;
let r1_result = r1.list_segments("").await;
assert!(
matches!(r1_result, Err(PagedbError::Aborted)),
"R1 should be aborted, got: {r1_result:?}"
);
let r2_result = r2.get(b"hello").await;
assert!(r2_result.is_ok(), "R2 should still work: {r2_result:?}");
drop(r1);
drop(r2);
}
#[tokio::test(flavor = "current_thread")]
async fn reopen_with_drainable_backlog_does_not_brick_readers() {
const THRESHOLD: u64 = 50;
let opts = || {
let mut o = OpenOptions::default()
.with_buffer_pool_pages(256)
.with_commit_history_retain(pagedb::options::RetainPolicy::Disabled);
o.reader_stall_threshold_pages = THRESHOLD;
o
};
let vfs = MemVfs::new();
{
let db = Db::open(vfs.clone(), KEK, 4096, REALM, opts())
.await
.unwrap();
{
let mut w = db.begin_write().await.unwrap();
w.put(b"seed", b"v").await.unwrap();
w.commit().await.unwrap();
}
let pin = db.begin_read().await.unwrap();
for round in 0u32..40 {
let mut w = db.begin_write().await.unwrap();
for i in 0u32..8 {
w.put(format!("k{round:03}_{i}").as_bytes(), &[0u8; 128])
.await
.unwrap();
}
w.commit().await.unwrap();
let mut w2 = db.begin_write().await.unwrap();
for i in 0u32..8 {
w2.put(format!("k{round:03}_{i}").as_bytes(), b"x")
.await
.unwrap();
}
w2.commit().await.unwrap();
}
let pending = db.stats().await.unwrap().free_list_pending_entries;
assert!(
pending > THRESHOLD,
"test setup should have built a backlog past the threshold; got {pending}"
);
drop(pin);
}
let db = Db::open(vfs.clone(), KEK, 4096, REALM, opts())
.await
.unwrap();
let reader = db.begin_read().await.unwrap();
let mut w = db.begin_write().await.unwrap();
w.put(b"schema", b"v1").await.unwrap();
let res = w.commit().await;
assert!(
res.is_ok(),
"post-reopen commit must not be rejected by the stall policy on an \
inherited drainable backlog: {res:?}"
);
let got = reader.get(b"seed").await;
assert!(
!matches!(got, Err(PagedbError::Aborted)),
"reopening a store with a drainable backlog bricked a live reader with \
Aborted — a non-corrupt, recoverable backlog must self-heal"
);
assert!(got.is_ok(), "reader read failed after reopen: {got:?}");
drop(reader);
}
#[tokio::test(flavor = "current_thread")]
async fn non_abortable_reader_survives_abort_oldest() {
let vfs = MemVfs::new();
let opts = opts_low_threshold();
let db = Db::open(vfs, KEK, 4096, REALM, opts).await.unwrap();
db.set_reader_stall_policy(ReaderStallPolicy::AbortOldest);
let r_non_abortable = db.begin_read_non_abortable().await.unwrap();
{
let mut txn = db.begin_write().await.unwrap();
txn.put(b"key", b"val").await.unwrap();
txn.commit().await.unwrap();
}
let err = fill_until_stall(&db).await;
assert!(
matches!(err, Some(PagedbError::DeferredFreeBacklog { .. })),
"should get DeferredFreeBacklog when only non-abortable readers are blocking; got {err:?}"
);
let result = r_non_abortable.get(b"key").await;
assert!(
result.is_ok(),
"non-abortable reader should not be aborted: {result:?}"
);
}