use crate::{AbstractTree, AnyTree, Config, SequenceNumberCounter, WriteBatch};
use std::sync::Arc;
use std::sync::mpsc;
#[test]
fn a_concurrent_write_cannot_land_while_the_scan_captures_the_memtable() {
let folder = tempfile::tempdir().expect("tempdir");
let any = Config::new(
folder.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()
.expect("open");
let AnyTree::Standard(tree) = any else {
panic!("expected a standard tree");
};
tree.insert("k", "v", 10);
let tree_for_hook = Arc::new(tree);
let writer_tree = Arc::clone(&tree_for_hook);
let (started_tx, started_rx) = mpsc::channel::<()>();
let (done_tx, done_rx) = mpsc::channel::<()>();
super::inner::TestHooks::install(
&tree_for_hook.test_hooks.scan_freeze,
Box::new(move || {
let writer_tree = Arc::clone(&writer_tree);
let done_tx = done_tx.clone();
let started_tx = started_tx.clone();
std::thread::spawn(move || {
started_tx.send(()).ok();
let mut batch = WriteBatch::new();
batch.insert("a", "concurrent");
batch.insert("z", "concurrent");
writer_tree.apply_batch(batch, 10).expect("apply");
done_tx.send(()).ok();
});
started_rx.recv().expect("writer thread started");
assert!(
done_rx
.recv_timeout(std::time::Duration::from_millis(250))
.is_err(),
"a writer must not be able to commit while the scan captures the \
active memtable — its entries would land mid-walk",
);
}),
);
let events: Vec<_> = tree_for_hook
.scan_since_seqno(0)
.expect("scan")
.collect::<Vec<_>>();
let keys: Vec<_> = events.iter().map(|e| e.key().to_vec()).collect::<Vec<_>>();
assert_eq!(
keys,
vec![b"k".to_vec()],
"the frozen capture holds exactly the state at scan start: {events:?}",
);
}
#[test]
fn a_range_deletion_holds_the_write_exclusion_guard_through_its_insert() {
use crate::ScanSinceEvent;
let folder = tempfile::tempdir().expect("tempdir");
let any = Config::new(
folder.path(),
SequenceNumberCounter::default(),
SequenceNumberCounter::default(),
)
.open()
.expect("open");
let AnyTree::Standard(tree) = any else {
panic!("expected a standard tree");
};
tree.insert("k", "v", 10);
let tree = Arc::new(tree);
let (started_tx, started_rx) = mpsc::channel::<()>();
let (gate_tx, gate_rx) = mpsc::channel::<()>();
super::inner::TestHooks::install(
&tree.test_hooks.range_write,
Box::new(move || {
started_tx.send(()).ok();
gate_rx.recv().expect("gate released");
}),
);
let writer_tree = Arc::clone(&tree);
let writer = std::thread::spawn(move || {
writer_tree.remove_range("a", "z", 10)
});
started_rx.recv().expect("writer reached its insert");
assert!(
tree.version_history.try_write().is_none(),
"a range-deletion writer must hold the version-history read guard \
through its tombstone insert — the CDC freeze relies on it",
);
gate_tx.send(()).expect("release writer");
writer.join().expect("writer thread");
let events: Vec<_> = tree.scan_since_seqno(0).expect("scan").collect::<Vec<_>>();
assert!(
events.iter().any(|e| matches!(
e,
ScanSinceEvent::RangeTombstone { start_key, end_key, seqno: 10 }
if start_key.as_ref() == b"a" && end_key.as_ref() == b"z"
)),
"the committed range deletion must surface as a CDC event: {events:?}",
);
}