use std::ops::{Deref, DerefMut};
use std::sync::{Arc, OnceLock};
use crate::sync::{Mutex, MutexGuard, RwLock};
use super::manifest::{Version, VersionSet};
use super::memtable::MemTable;
pub(crate) struct ReadView {
pub(crate) active: Arc<MemTable>,
pub(crate) frozen: Vec<Arc<MemTable>>,
pub(crate) version: Arc<Version>,
}
pub(crate) struct ReadViewCell {
current: RwLock<Arc<ReadView>>,
publish: Mutex<()>,
}
impl ReadViewCell {
pub(crate) fn new(view: ReadView) -> Self {
Self {
current: RwLock::new(Arc::new(view)),
publish: Mutex::new(()),
}
}
#[inline]
pub(crate) fn load(&self) -> Arc<ReadView> {
Arc::clone(&self.current.read())
}
pub(crate) fn update_memtables<R>(
&self,
mutate: impl FnOnce(&Arc<MemTable>, &[Arc<MemTable>]) -> (Arc<MemTable>, Vec<Arc<MemTable>>, R),
) -> R {
let _publishing = self.publish.lock();
let current = self.load();
let (active, frozen, out) = mutate(¤t.active, ¤t.frozen);
let next = Arc::new(ReadView {
active,
frozen,
version: Arc::clone(¤t.version),
});
*self.current.write() = next;
out
}
pub(crate) fn retire_memtable(&self, flushed: &Arc<MemTable>) {
self.update_memtables(|active, frozen| {
let next = frozen
.iter()
.filter(|mt| !Arc::ptr_eq(mt, flushed))
.cloned()
.collect();
(Arc::clone(active), next, ())
});
}
fn publish_version(&self, version: Arc<Version>) {
let _publishing = self.publish.lock();
let current = self.load();
let next = Arc::new(ReadView {
active: Arc::clone(¤t.active),
frozen: current.frozen.clone(),
version,
});
*self.current.write() = next;
}
}
pub(crate) struct VersionStore {
inner: Mutex<VersionSet>,
view: OnceLock<Arc<ReadViewCell>>,
}
impl VersionStore {
pub(crate) fn new(versions: VersionSet) -> Self {
Self {
inner: Mutex::new(versions),
view: OnceLock::new(),
}
}
pub(crate) fn attach_view(&self, cell: Arc<ReadViewCell>) {
let _ = self.view.set(cell);
}
pub(crate) fn lock(&self) -> VersionGuard<'_> {
let guard = self.inner.lock();
let entry_version = guard.current();
VersionGuard {
entry_version,
guard,
view: self.view.get(),
}
}
}
pub(crate) struct VersionGuard<'a> {
guard: MutexGuard<'a, VersionSet>,
view: Option<&'a Arc<ReadViewCell>>,
entry_version: Arc<Version>,
}
impl Deref for VersionGuard<'_> {
type Target = VersionSet;
fn deref(&self) -> &VersionSet {
&self.guard
}
}
impl DerefMut for VersionGuard<'_> {
fn deref_mut(&mut self) -> &mut VersionSet {
&mut self.guard
}
}
impl Drop for VersionGuard<'_> {
fn drop(&mut self) {
let Some(view) = self.view else {
return;
};
let current = self.guard.current();
if Arc::ptr_eq(¤t, &self.entry_version) {
return;
}
view.publish_version(current);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::engine::manifest::VersionEdit;
fn test_memtable_config() -> crate::engine::memtable::MemTableConfig {
crate::engine::memtable::MemTableConfig::new(
crate::engine::arena::ArenaProfile::EMBEDDED,
64 * 1024,
2,
)
}
fn store_with_view() -> (tempfile::TempDir, Arc<VersionStore>, Arc<ReadViewCell>) {
let dir = tempfile::tempdir().unwrap();
let sst_dir = dir.path().join("sst");
std::fs::create_dir_all(&sst_dir).unwrap();
let versions = VersionSet::open(dir.path(), &sst_dir).unwrap();
let store = Arc::new(VersionStore::new(versions));
let cell = Arc::new(ReadViewCell::new(ReadView {
active: Arc::new(MemTable::new(&test_memtable_config()).unwrap()),
frozen: Vec::new(),
version: store.lock().current(),
}));
store.attach_view(Arc::clone(&cell));
(dir, store, cell)
}
#[test]
fn retiring_a_memtable_drops_that_one_and_leaves_the_rest_in_order() {
let (_dir, _store, cell) = store_with_view();
let frozen: Vec<Arc<MemTable>> = (0..3)
.map(|i| {
let mt = Arc::new(MemTable::new(&test_memtable_config()).unwrap());
mt.put(
format!("k{i}").as_bytes(),
format!("v{i}").as_bytes(),
i + 1,
);
mt
})
.collect();
cell.update_memtables(|active, _| (Arc::clone(active), frozen.clone(), ()));
cell.retire_memtable(&frozen[1]);
let after = cell.load();
assert_eq!(after.frozen.len(), 2, "exactly one memtable must go");
assert!(
Arc::ptr_eq(&after.frozen[0], &frozen[0]),
"retiring the middle memtable dropped the oldest one instead: every write in it is \
in no published version and is now unreachable",
);
assert!(
Arc::ptr_eq(&after.frozen[1], &frozen[2]),
"the newest memtable did not keep its place",
);
}
#[test]
fn a_sealed_memtable_carries_the_log_its_records_are_in() {
let (_dir, _store, cell) = store_with_view();
let first = Arc::new(MemTable::new(&test_memtable_config()).unwrap());
first.put(b"a", b"1", 1);
first.seal_wal(std::path::PathBuf::from("/wal/000001.log"));
let second = Arc::new(MemTable::new(&test_memtable_config()).unwrap());
second.put(b"b", b"2", 2);
second.seal_wal(std::path::PathBuf::from("/wal/000002.log"));
cell.update_memtables(|active, _| {
(Arc::clone(active), vec![first.clone(), second.clone()], ())
});
let view = cell.load();
assert_eq!(
view.frozen[0].sealed_wal(),
Some(std::path::Path::new("/wal/000001.log")),
"the front of the frozen list must name its own log, not the newest one",
);
assert_eq!(
view.frozen[1].sealed_wal(),
Some(std::path::Path::new("/wal/000002.log")),
);
assert_eq!(
view.active.sealed_wal(),
None,
"the active memtable is still taking writes, so its log is not sealed",
);
cell.retire_memtable(&first);
let after = cell.load();
assert_eq!(after.frozen.len(), 1);
assert_eq!(
after.frozen[0].sealed_wal(),
Some(std::path::Path::new("/wal/000002.log")),
);
}
#[test]
fn retiring_a_memtable_that_is_already_gone_changes_nothing() {
let (_dir, _store, cell) = store_with_view();
let frozen: Vec<Arc<MemTable>> = (0..2)
.map(|_| Arc::new(MemTable::new(&test_memtable_config()).unwrap()))
.collect();
cell.update_memtables(|active, _| (Arc::clone(active), frozen.clone(), ()));
cell.retire_memtable(&frozen[0]);
cell.retire_memtable(&frozen[0]);
let after = cell.load();
assert_eq!(
after.frozen.len(),
1,
"a second retirement of the same memtable took a different one with it",
);
assert!(Arc::ptr_eq(&after.frozen[0], &frozen[1]));
}
#[test]
fn a_rotation_publishes_the_sealed_memtable_and_the_fresh_one_together() {
let (_dir, _store, cell) = store_with_view();
let before = cell.load();
before.active.put(b"k", b"v", 1);
cell.update_memtables(|active, frozen| {
let mut next = frozen.to_vec();
next.push(Arc::clone(active));
(
Arc::new(MemTable::new(&test_memtable_config()).unwrap()),
next,
(),
)
});
let after = cell.load();
assert!(after.active.is_empty(), "writers got a fresh memtable");
assert_eq!(after.frozen.len(), 1);
assert!(
Arc::ptr_eq(&after.frozen[0], &before.active),
"the sealed memtable is the one writers were using",
);
assert!(
Arc::ptr_eq(&after.version, &before.version),
"a memtable publication leaves the version alone",
);
assert!(
!before.active.is_empty(),
"the view a reader still holds keeps its data",
);
}
#[test]
fn a_version_edit_publishes_a_new_view_that_keeps_the_memtables() {
let (_dir, store, cell) = store_with_view();
let before = cell.load();
store
.lock()
.apply(&[VersionEdit::SetNextFileId(7)])
.unwrap();
let after = cell.load();
assert_eq!(after.version.next_file_id, 7);
assert!(Arc::ptr_eq(&after.active, &before.active));
assert_eq!(after.frozen.len(), before.frozen.len());
}
#[test]
fn a_critical_section_that_changes_no_version_publishes_nothing() {
let (_dir, store, cell) = store_with_view();
let before = cell.load();
{
let guard = store.lock();
let _ = guard.current();
}
assert!(
Arc::ptr_eq(&before, &cell.load()),
"a read-only critical section must not churn the published view",
);
}
}