use std::io::{self};
use std::path::{Path, PathBuf};
#[cfg(test)]
use std::fs::OpenOptions;
use std::sync::Arc;
use crate::env::{BufferedWriter, Env, WriteMode};
use crate::sync::RwLock;
use super::checksum;
use super::sstable::{
LiveSst, MetadataPolicy, SsTableMeta, SsTableReader, sst_filename, table_carries_data,
};
pub(crate) const MAX_LEVELS: usize = 7;
#[derive(Clone)]
pub(crate) struct Version {
pub(crate) levels: Vec<Vec<Arc<LiveSst>>>,
pub(crate) next_file_id: u64,
pub(crate) last_seq: u64,
pub(crate) min_wal_id: u64,
}
impl Version {
pub(crate) fn new() -> Self {
Self {
levels: (0..MAX_LEVELS).map(|_| Vec::new()).collect(),
next_file_id: 1,
last_seq: 0,
min_wal_id: 0,
}
}
pub(crate) fn l0_count(&self) -> usize {
self.levels[0].len()
}
pub(crate) fn level_size(&self, level: usize) -> u64 {
self.levels[level].iter().map(|f| f.meta.file_size).sum()
}
}
#[derive(Clone)]
pub(crate) enum VersionEdit {
AddFile { level: usize, file: Arc<LiveSst> },
RemoveFile { level: usize, file_id: u64 },
SetLastSeq(u64),
SetNextFileId(u64),
Reset { next_file_id: u64, min_wal_id: u64 },
}
enum ManifestRecord {
AddFile { level: usize, meta: SsTableMeta },
RemoveFile { level: usize, file_id: u64 },
SetLastSeq(u64),
SetNextFileId(u64),
SetMinWalId(u64),
Reset { next_file_id: u64, min_wal_id: u64 },
}
const TAG_ADD_FILE: u8 = 1;
const TAG_REMOVE_FILE: u8 = 2;
const TAG_LAST_SEQ: u8 = 3;
const TAG_NEXT_FILE_ID: u8 = 4;
const TAG_MIN_WAL_ID: u8 = 5;
const TAG_RESET: u8 = 6;
impl VersionEdit {
fn to_record(&self) -> ManifestRecord {
match self {
VersionEdit::AddFile { level, file } => ManifestRecord::AddFile {
level: *level,
meta: file.meta.clone(),
},
VersionEdit::RemoveFile { level, file_id } => ManifestRecord::RemoveFile {
level: *level,
file_id: *file_id,
},
VersionEdit::SetLastSeq(seq) => ManifestRecord::SetLastSeq(*seq),
VersionEdit::SetNextFileId(id) => ManifestRecord::SetNextFileId(*id),
VersionEdit::Reset {
next_file_id,
min_wal_id,
} => ManifestRecord::Reset {
next_file_id: *next_file_id,
min_wal_id: *min_wal_id,
},
}
}
fn requires_manifest_sync(&self) -> bool {
!matches!(self, VersionEdit::SetNextFileId(_))
}
}
impl ManifestRecord {
fn encode(&self, buf: &mut Vec<u8>) {
match self {
ManifestRecord::AddFile { level, meta } => {
buf.push(TAG_ADD_FILE);
buf.extend_from_slice(&(*level as u32).to_le_bytes());
buf.extend_from_slice(&meta.file_id.to_le_bytes());
buf.extend_from_slice(&(meta.smallest_key.len() as u32).to_le_bytes());
buf.extend_from_slice(&meta.smallest_key);
buf.extend_from_slice(&(meta.largest_key.len() as u32).to_le_bytes());
buf.extend_from_slice(&meta.largest_key);
buf.extend_from_slice(&meta.file_size.to_le_bytes());
buf.extend_from_slice(&meta.num_entries.to_le_bytes());
}
ManifestRecord::RemoveFile { level, file_id } => {
buf.push(TAG_REMOVE_FILE);
buf.extend_from_slice(&(*level as u32).to_le_bytes());
buf.extend_from_slice(&file_id.to_le_bytes());
}
ManifestRecord::SetLastSeq(seq) => {
buf.push(TAG_LAST_SEQ);
buf.extend_from_slice(&seq.to_le_bytes());
}
ManifestRecord::SetNextFileId(id) => {
buf.push(TAG_NEXT_FILE_ID);
buf.extend_from_slice(&id.to_le_bytes());
}
ManifestRecord::SetMinWalId(id) => {
buf.push(TAG_MIN_WAL_ID);
buf.extend_from_slice(&id.to_le_bytes());
}
ManifestRecord::Reset {
next_file_id,
min_wal_id,
} => {
buf.push(TAG_RESET);
buf.extend_from_slice(&next_file_id.to_le_bytes());
buf.extend_from_slice(&min_wal_id.to_le_bytes());
}
}
}
fn decode(data: &[u8], pos: &mut usize) -> io::Result<Option<Self>> {
if *pos >= data.len() {
return Ok(None);
}
let tag = data[*pos];
*pos += 1;
match tag {
TAG_ADD_FILE => {
let level = read_u32(data, pos)? as usize;
validate_level_index(level)?;
let file_id = read_u64(data, pos)?;
let smallest_key = read_bytes(data, pos)?;
let largest_key = read_bytes(data, pos)?;
let file_size = read_u64(data, pos)?;
let num_entries = read_u64(data, pos)?;
Ok(Some(ManifestRecord::AddFile {
level,
meta: SsTableMeta {
file_id,
smallest_key,
largest_key,
file_size,
num_entries,
},
}))
}
TAG_REMOVE_FILE => {
let level = read_u32(data, pos)? as usize;
validate_level_index(level)?;
let file_id = read_u64(data, pos)?;
Ok(Some(ManifestRecord::RemoveFile { level, file_id }))
}
TAG_LAST_SEQ => {
let seq = read_u64(data, pos)?;
Ok(Some(ManifestRecord::SetLastSeq(seq)))
}
TAG_NEXT_FILE_ID => {
let id = read_u64(data, pos)?;
Ok(Some(ManifestRecord::SetNextFileId(id)))
}
TAG_MIN_WAL_ID => {
let id = read_u64(data, pos)?;
Ok(Some(ManifestRecord::SetMinWalId(id)))
}
TAG_RESET => {
let next_file_id = read_u64(data, pos)?;
let min_wal_id = read_u64(data, pos)?;
Ok(Some(ManifestRecord::Reset {
next_file_id,
min_wal_id,
}))
}
_ => Err(io::Error::new(
io::ErrorKind::InvalidData,
format!("unknown manifest record tag: {}", tag),
)),
}
}
}
fn read_u32(data: &[u8], pos: &mut usize) -> io::Result<u32> {
if *pos + 4 > data.len() {
return Err(io::Error::new(io::ErrorKind::UnexpectedEof, "short read"));
}
let val = u32::from_le_bytes(data[*pos..*pos + 4].try_into().unwrap());
*pos += 4;
Ok(val)
}
fn read_u64(data: &[u8], pos: &mut usize) -> io::Result<u64> {
if *pos + 8 > data.len() {
return Err(io::Error::new(io::ErrorKind::UnexpectedEof, "short read"));
}
let val = u64::from_le_bytes(data[*pos..*pos + 8].try_into().unwrap());
*pos += 8;
Ok(val)
}
fn read_bytes(data: &[u8], pos: &mut usize) -> io::Result<Vec<u8>> {
let len = read_u32(data, pos)? as usize;
if *pos + len > data.len() {
return Err(io::Error::new(io::ErrorKind::UnexpectedEof, "short read"));
}
let bytes = data[*pos..*pos + len].to_vec();
*pos += len;
Ok(bytes)
}
fn validate_level_index(level: usize) -> io::Result<()> {
if level < MAX_LEVELS {
return Ok(());
}
Err(io::Error::new(
io::ErrorKind::InvalidData,
format!("manifest level {level} out of range 0..{MAX_LEVELS}"),
))
}
const MANIFEST_MAGIC: [u8; 7] = *b"REGOMAN";
const MANIFEST_FORMAT_V1: u8 = 1;
const MANIFEST_STAMP_LEN: usize = 12;
const MANIFEST_COMPACT_FACTOR: u64 = 4;
const MANIFEST_COMPACT_FLOOR: u64 = 64 * 1024;
const APPROX_ADD_FILE_RECORD_BYTES: u64 = 256;
pub(crate) struct VersionSet {
current: Arc<RwLock<Arc<Version>>>,
manifest_path: PathBuf,
manifest_bytes: u64,
manifest_writer: Option<BufferedWriter>,
env: Arc<dyn Env>,
}
struct ManifestReplay {
version: Version,
valid_len: usize,
}
struct SuspectTable {
path: PathBuf,
len: Option<u64>,
reason: String,
}
const SUSPECTS_NAMED: usize = 8;
fn describe_suspects(suspects: &[SuspectTable]) -> String {
let mut out = suspects
.iter()
.take(SUSPECTS_NAMED)
.map(|s| match s.len {
Some(len) => format!("{} ({len} bytes, {})", s.path.display(), s.reason),
None => format!("{} (size unknown, {})", s.path.display(), s.reason),
})
.collect::<Vec<_>>()
.join(", ");
if suspects.len() > SUSPECTS_NAMED {
out.push_str(&format!(
", and {} more not named here",
suspects.len() - SUSPECTS_NAMED
));
}
out
}
impl VersionSet {
pub(crate) fn open_with_policy(
env: &Arc<dyn Env>,
db_dir: &Path,
sst_dir: &Path,
policy: MetadataPolicy,
) -> io::Result<Self> {
let manifest_path = db_dir.join("MANIFEST");
let manifest_bytes;
let (version, writer) = if env.exists(&manifest_path) {
let data = env.read(&manifest_path)?;
let replay = Self::replay_manifest(env, &data, sst_dir, policy)?;
Self::reject_discarded_tables(&**env, &replay, data.len(), sst_dir, &manifest_path)?;
if replay.valid_len < data.len() {
let mut trim = env.open_write(&manifest_path, WriteMode::Update)?;
trim.set_len(replay.valid_len as u64)?;
trim.sync_all()?;
}
if data.len() < MANIFEST_STAMP_LEN {
let mut file = env.open_write(&manifest_path, WriteMode::Truncate)?;
file.write_all(&Self::encode_stamp())?;
file.sync_all()?;
crate::env::sync_parent_dir(&**env, &manifest_path)?;
manifest_bytes = MANIFEST_STAMP_LEN as u64;
(replay.version, BufferedWriter::new(file))
} else {
let file = env.open_write(&manifest_path, WriteMode::Append)?;
manifest_bytes = replay.valid_len as u64;
(replay.version, BufferedWriter::new(file))
}
} else {
let version = Version::new();
let mut file = env.open_write(&manifest_path, WriteMode::Truncate)?;
file.write_all(&Self::encode_stamp())?;
file.sync_all()?;
crate::env::sync_parent_dir(&**env, &manifest_path)?;
manifest_bytes = MANIFEST_STAMP_LEN as u64;
(version, BufferedWriter::new(file))
};
Ok(Self {
current: Arc::new(RwLock::new(Arc::new(version))),
manifest_path,
manifest_bytes,
manifest_writer: Some(writer),
env: Arc::clone(env),
})
}
#[cfg(any(test, feature = "fuzzing"))]
pub(crate) fn open(db_dir: &Path, sst_dir: &Path) -> io::Result<Self> {
Self::open_with_policy(
&crate::env::std_env(),
db_dir,
sst_dir,
MetadataPolicy::Pinned,
)
}
pub(crate) fn open_read_only(
env: &Arc<dyn Env>,
db_dir: &Path,
sst_dir: &Path,
policy: MetadataPolicy,
) -> io::Result<Self> {
let manifest_path = db_dir.join("MANIFEST");
let data = env.read(&manifest_path)?;
let replay = Self::replay_manifest(env, &data, sst_dir, policy)?;
Self::reject_discarded_tables(&**env, &replay, data.len(), sst_dir, &manifest_path)?;
Ok(Self {
current: Arc::new(RwLock::new(Arc::new(replay.version))),
manifest_path,
manifest_bytes: replay.valid_len as u64,
manifest_writer: None,
env: Arc::clone(env),
})
}
pub(crate) fn current(&self) -> Arc<Version> {
Arc::clone(&*self.current.read())
}
pub(crate) fn manifest_path(&self) -> &Path {
&self.manifest_path
}
pub(crate) fn apply(&mut self, edits: &[VersionEdit]) -> io::Result<()> {
for edit in edits {
match edit {
VersionEdit::AddFile { level, .. } | VersionEdit::RemoveFile { level, .. } => {
validate_level_index(*level)?;
}
VersionEdit::SetLastSeq(_)
| VersionEdit::SetNextFileId(_)
| VersionEdit::Reset { .. } => {}
}
}
let mut version = (*self.current()).clone();
for edit in edits {
match edit {
VersionEdit::AddFile { level, file } => {
version.levels[*level].push(Arc::clone(file));
}
VersionEdit::RemoveFile { level, file_id } => {
version.levels[*level].retain(|f| f.meta.file_id != *file_id);
}
VersionEdit::SetLastSeq(seq) => {
version.last_seq = *seq;
}
VersionEdit::SetNextFileId(id) => {
version.next_file_id = *id;
}
VersionEdit::Reset {
next_file_id,
min_wal_id,
} => {
version.levels = (0..MAX_LEVELS).map(|_| Vec::new()).collect();
version.last_seq = 0;
version.next_file_id = *next_file_id;
version.min_wal_id = *min_wal_id;
}
}
}
let records: Vec<ManifestRecord> = edits.iter().map(VersionEdit::to_record).collect();
let encoded = Self::encode_records(&records);
let requires_sync = edits.iter().any(VersionEdit::requires_manifest_sync);
if let Some(writer) = &mut self.manifest_writer {
writer.write_all(&encoded)?;
if requires_sync {
writer.sync_all()?;
} else {
writer.flush()?;
}
}
let live_files: u64 = version.levels.iter().map(|l| l.len() as u64).sum();
*self.current.write() = Arc::new(version);
self.manifest_bytes += encoded.len() as u64;
if self.manifest_bytes > Self::compact_threshold(live_files) {
self.compact_manifest()?;
}
Ok(())
}
fn compact_threshold(live_files: u64) -> u64 {
let canonical = MANIFEST_STAMP_LEN as u64 + live_files * APPROX_ADD_FILE_RECORD_BYTES;
MANIFEST_COMPACT_FLOOR.max(canonical.saturating_mul(MANIFEST_COMPACT_FACTOR))
}
pub(crate) fn compact_manifest(&mut self) -> io::Result<()> {
let version = self.current();
let mut records = Vec::new();
records.push(ManifestRecord::SetNextFileId(version.next_file_id));
records.push(ManifestRecord::SetLastSeq(version.last_seq));
records.push(ManifestRecord::SetMinWalId(version.min_wal_id));
for (level, files) in version.levels.iter().enumerate() {
for file in files {
records.push(ManifestRecord::AddFile {
level,
meta: file.meta.clone(),
});
}
}
let encoded = Self::encode_records(&records);
let tmp_path = self.manifest_path.with_extension("tmp");
{
let mut file = self.env.open_write(&tmp_path, WriteMode::Truncate)?;
file.write_all(&Self::encode_stamp())?;
file.write_all(&encoded)?;
file.sync_all()?;
}
self.manifest_writer = None;
self.env.rename(&tmp_path, &self.manifest_path)?;
crate::env::sync_parent_dir(&*self.env, &self.manifest_path)?;
self.manifest_bytes = (MANIFEST_STAMP_LEN + encoded.len()) as u64;
let file = self
.env
.open_write(&self.manifest_path, WriteMode::Append)?;
self.manifest_writer = Some(BufferedWriter::new(file));
Ok(())
}
pub(crate) fn encode_stamp() -> [u8; MANIFEST_STAMP_LEN] {
let mut out = [0u8; MANIFEST_STAMP_LEN];
out[0..7].copy_from_slice(&MANIFEST_MAGIC);
out[7] = MANIFEST_FORMAT_V1;
let checksum = checksum::manifest_record(0, &out[0..8]);
out[8..12].copy_from_slice(&checksum.to_le_bytes());
out
}
fn stamp_len(data: &[u8]) -> io::Result<usize> {
if data.len() < MANIFEST_STAMP_LEN {
return Ok(0);
}
if data[0..7] != MANIFEST_MAGIC {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"MANIFEST does not begin with the REGOMAN stamp",
));
}
let format = data[7];
let stored = u32::from_le_bytes([data[8], data[9], data[10], data[11]]);
if stored != checksum::manifest_record(0, &data[0..8]) {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"MANIFEST stamp checksum mismatch",
));
}
if format > MANIFEST_FORMAT_V1 {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
format!(
"MANIFEST format {format} was written by a newer regolith than this build, \
which understands up to {MANIFEST_FORMAT_V1}"
),
));
}
Ok(MANIFEST_STAMP_LEN)
}
fn encode_records(records: &[ManifestRecord]) -> Vec<u8> {
let mut buf = Vec::new();
for record in records {
let mut record_buf = Vec::new();
record.encode(&mut record_buf);
let len = record_buf.len() as u32;
let checksum = checksum::manifest_record(len, &record_buf);
buf.extend_from_slice(&len.to_le_bytes());
buf.extend_from_slice(&record_buf);
buf.extend_from_slice(&checksum.to_le_bytes());
}
buf
}
fn reject_discarded_tables(
env: &dyn Env,
replay: &ManifestReplay,
manifest_len: usize,
sst_dir: &Path,
manifest_path: &Path,
) -> io::Result<()> {
let replayed_cleanly = manifest_len > 0 && replay.valid_len == manifest_len;
if replayed_cleanly || replay.version.levels.iter().any(|level| !level.is_empty()) {
return Ok(());
}
let suspects = Self::suspect_tables(env, sst_dir);
if suspects.is_empty() {
return Ok(());
}
Err(io::Error::new(
io::ErrorKind::InvalidData,
format!(
"{} is corrupt: it references no SSTable, but {} table file(s) in {} may still hold data. \
Opening would discard them, so the database is left untouched. Suspect tables: {}",
manifest_path.display(),
suspects.len(),
sst_dir.display(),
describe_suspects(&suspects),
),
))
}
fn suspect_tables(env: &dyn Env, sst_dir: &Path) -> Vec<SuspectTable> {
let entries = match env.read_dir(sst_dir) {
Ok(entries) => entries,
Err(e) => {
tracing::warn!(
dir = %sst_dir.display(),
error = %e,
"could not list the SSTable directory while checking for discarded tables"
);
return Vec::new();
}
};
let mut suspects = Vec::new();
for entry in entries {
let path = entry.path.clone();
if path.extension().and_then(|ext| ext.to_str()) != Some("sst") {
continue;
}
let len = match env.metadata(&path) {
Ok(meta) => Some(meta.len),
Err(e) => {
suspects.push(SuspectTable {
path,
len: None,
reason: format!("unreadable: {e}"),
});
continue;
}
};
if len == Some(0) {
tracing::warn!(
path = %path.display(),
"ignoring a zero-length orphan SSTable left by a crash inside a flush"
);
continue;
}
match table_carries_data(env, &path) {
Ok(true) => suspects.push(SuspectTable {
path,
len,
reason: "carries data".to_string(),
}),
Ok(false) => tracing::warn!(
path = %path.display(),
"ignoring an orphan SSTable whose footer records no entry and no range tombstone"
),
Err(e) => suspects.push(SuspectTable {
path,
len,
reason: format!("unreadable footer: {e}"),
}),
}
}
suspects.sort_by(|a, b| a.path.cmp(&b.path));
suspects
}
fn replay_manifest(
env: &Arc<dyn Env>,
data: &[u8],
sst_dir: &Path,
policy: MetadataPolicy,
) -> io::Result<ManifestReplay> {
let mut surviving: Vec<Vec<SsTableMeta>> = vec![Vec::new(); MAX_LEVELS];
let mut last_seq: u64 = 0;
let mut next_file_id: u64 = 1;
let mut min_wal_id: u64 = 0;
let stamp = Self::stamp_len(data)?;
let mut offset = stamp;
let mut valid_len = stamp;
while offset < data.len() {
if offset + 4 > data.len() {
tracing::warn!("Truncated manifest record header, stopping replay");
break;
}
let len = u32::from_le_bytes(data[offset..offset + 4].try_into().unwrap()) as usize;
offset += 4;
if offset + len + 4 > data.len() {
tracing::warn!("Truncated manifest record, stopping replay");
break;
}
let record_data = &data[offset..offset + len];
offset += len;
let stored_checksum = u32::from_le_bytes(data[offset..offset + 4].try_into().unwrap());
offset += 4;
let computed_checksum = checksum::manifest_record(len as u32, record_data);
if stored_checksum != computed_checksum {
tracing::warn!("Manifest checksum mismatch, stopping replay");
break;
}
let mut pos = 0;
while let Some(record) = ManifestRecord::decode(record_data, &mut pos)? {
match record {
ManifestRecord::AddFile { level, meta } => {
surviving[level].push(meta);
}
ManifestRecord::RemoveFile { level, file_id } => {
surviving[level].retain(|m| m.file_id != file_id);
}
ManifestRecord::SetLastSeq(seq) => {
last_seq = seq;
}
ManifestRecord::SetNextFileId(id) => {
next_file_id = id;
}
ManifestRecord::SetMinWalId(id) => {
min_wal_id = id;
}
ManifestRecord::Reset {
next_file_id: reset_next_file_id,
min_wal_id: reset_min_wal_id,
} => {
for level in &mut surviving {
level.clear();
}
last_seq = 0;
next_file_id = reset_next_file_id;
min_wal_id = reset_min_wal_id;
}
}
}
valid_len = offset;
}
let mut version = Version::new();
version.last_seq = last_seq;
version.next_file_id = next_file_id;
version.min_wal_id = min_wal_id;
for (level, files) in surviving.into_iter().enumerate() {
for meta in files {
let path = sst_dir.join(sst_filename(meta.file_id));
let reader = Arc::new(
SsTableReader::open_with(env, &path, meta.file_id, policy).map_err(|e| {
std::io::Error::new(e.kind(), format!("open {}: {e}", path.display()))
})?,
);
version.levels[level].push(LiveSst::new(meta, reader));
}
}
Ok(ManifestReplay { version, valid_len })
}
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
fn make_live_sst(dir: &Path, file_id: u64, smallest: &[u8], largest: &[u8]) -> Arc<LiveSst> {
use super::super::internal_key::{VALUE_TYPE_VALUE, encode_internal_key};
use super::super::sstable::SsTableWriter;
use crate::options::CompressionType;
let path = dir.join(sst_filename(file_id));
let mut writer =
SsTableWriter::new(&path, 4096, 10, CompressionType::None, None, false, 4096).unwrap();
writer
.add(
&encode_internal_key(smallest, 1, VALUE_TYPE_VALUE),
b"value",
)
.unwrap();
if smallest != largest {
writer
.add(&encode_internal_key(largest, 1, VALUE_TYPE_VALUE), b"value")
.unwrap();
}
let summary = writer.finish().unwrap().unwrap();
let file_size = std::fs::metadata(&path).unwrap().len();
let reader = Arc::new(SsTableReader::open(&path, file_id).unwrap());
LiveSst::new(
SsTableMeta {
file_id,
smallest_key: summary.smallest_user_key,
largest_key: summary.largest_user_key,
file_size,
num_entries: summary.num_entries,
},
reader,
)
}
fn second_record_checksum_offset(path: &Path) -> usize {
let data = std::fs::read(path).unwrap();
let base = MANIFEST_STAMP_LEN;
let first_len = u32::from_le_bytes(data[base..base + 4].try_into().unwrap()) as usize;
let second_start = base + 4 + first_len + 4;
let second_len =
u32::from_le_bytes(data[second_start..second_start + 4].try_into().unwrap()) as usize;
second_start + 4 + second_len
}
fn test_meta(file_id: u64) -> SsTableMeta {
SsTableMeta {
file_id,
smallest_key: b"a".to_vec(),
largest_key: b"z".to_vec(),
file_size: 128,
num_entries: 2,
}
}
#[test]
fn manifest_checksum_covers_length_header() {
let mut record = Vec::new();
ManifestRecord::SetLastSeq(7).encode(&mut record);
let len = record.len() as u32;
let baseline = checksum::manifest_record(len, &record);
assert_ne!(baseline, checksum::manifest_record(len + 1, &record));
}
#[test]
fn test_apply_and_replay_roundtrip() {
let dir = TempDir::new().unwrap();
let sst_dir = dir.path().join("sst");
std::fs::create_dir_all(&sst_dir).unwrap();
let file1 = make_live_sst(&sst_dir, 1, b"aaa", b"zzz");
let edits = vec![
VersionEdit::AddFile {
level: 0,
file: Arc::clone(&file1),
},
VersionEdit::SetLastSeq(42),
VersionEdit::SetNextFileId(10),
];
{
let mut vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
vs.apply(&edits).unwrap();
let v = vs.current();
assert_eq!(v.levels[0].len(), 1);
assert_eq!(v.levels[0][0].meta.file_id, 1);
assert_eq!(v.last_seq, 42);
assert_eq!(v.next_file_id, 10);
assert_eq!(v.min_wal_id, 0);
}
let vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
let v = vs.current();
assert_eq!(v.levels[0].len(), 1);
assert_eq!(v.levels[0][0].meta.file_id, 1);
assert_eq!(v.last_seq, 42);
assert_eq!(v.next_file_id, 10);
assert_eq!(v.min_wal_id, 0);
}
#[test]
fn test_remove_file_hides_it_from_new_version() {
let dir = TempDir::new().unwrap();
let sst_dir = dir.path().join("sst");
std::fs::create_dir_all(&sst_dir).unwrap();
let file1 = make_live_sst(&sst_dir, 1, b"a", b"c");
let file2 = make_live_sst(&sst_dir, 2, b"d", b"f");
let mut vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
vs.apply(&[
VersionEdit::AddFile {
level: 0,
file: Arc::clone(&file1),
},
VersionEdit::AddFile {
level: 0,
file: Arc::clone(&file2),
},
])
.unwrap();
let pinned = vs.current();
assert_eq!(pinned.levels[0].len(), 2);
vs.apply(&[VersionEdit::RemoveFile {
level: 0,
file_id: 1,
}])
.unwrap();
let v = vs.current();
assert_eq!(v.levels[0].len(), 1);
assert_eq!(v.levels[0][0].meta.file_id, 2);
assert_eq!(pinned.levels[0].len(), 2);
}
#[test]
fn initial_version_has_empty_levels_and_defaults() {
let v = Version::new();
assert_eq!(v.levels.len(), MAX_LEVELS);
assert!(v.levels.iter().all(|l| l.is_empty()));
assert_eq!(v.next_file_id, 1);
assert_eq!(v.last_seq, 0);
assert_eq!(v.min_wal_id, 0);
assert_eq!(v.l0_count(), 0);
for level in 0..MAX_LEVELS {
assert_eq!(v.level_size(level), 0);
}
}
#[test]
fn manifest_path_returned_as_written() {
let dir = TempDir::new().unwrap();
let sst_dir = dir.path().join("sst");
std::fs::create_dir_all(&sst_dir).unwrap();
let vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
assert_eq!(vs.manifest_path(), dir.path().join("MANIFEST"));
}
#[test]
fn manifest_record_encode_decode_round_trip() {
let records = [
ManifestRecord::AddFile {
level: 2,
meta: SsTableMeta {
file_id: 99,
smallest_key: b"aaa".to_vec(),
largest_key: b"zzz".to_vec(),
file_size: 4096,
num_entries: 128,
},
},
ManifestRecord::RemoveFile {
level: 1,
file_id: 7,
},
ManifestRecord::SetLastSeq(999),
ManifestRecord::SetNextFileId(42),
ManifestRecord::SetMinWalId(11),
ManifestRecord::Reset {
next_file_id: 77,
min_wal_id: 76,
},
];
for r in &records {
let mut buf = Vec::new();
r.encode(&mut buf);
let mut pos = 0;
let decoded = match ManifestRecord::decode(&buf, &mut pos) {
Ok(Some(d)) => d,
other => panic!("expected decoded record, got {:?}", other.is_ok()),
};
let mut rebuf = Vec::new();
decoded.encode(&mut rebuf);
assert_eq!(buf, rebuf);
assert_eq!(pos, buf.len());
}
}
#[test]
fn manifest_record_decode_rejects_unknown_tag() {
let data = [0xFFu8];
let mut pos = 0;
let kind = match ManifestRecord::decode(&data, &mut pos) {
Err(e) => e.kind(),
Ok(_) => panic!("expected error on unknown tag"),
};
assert_eq!(kind, io::ErrorKind::InvalidData);
}
#[test]
fn manifest_record_decode_rejects_invalid_level_indexes() {
let records = [
ManifestRecord::AddFile {
level: MAX_LEVELS,
meta: test_meta(1),
},
ManifestRecord::RemoveFile {
level: MAX_LEVELS,
file_id: 1,
},
];
for record in records {
let mut data = Vec::new();
record.encode(&mut data);
let mut pos = 0;
let kind = match ManifestRecord::decode(&data, &mut pos) {
Err(e) => e.kind(),
Ok(_) => panic!("expected invalid level error"),
};
assert_eq!(kind, io::ErrorKind::InvalidData);
}
}
#[test]
fn manifest_record_decode_returns_none_at_eof() {
let mut pos = 0;
let got = ManifestRecord::decode(&[], &mut pos).unwrap();
assert!(got.is_none());
}
#[test]
fn apply_rejects_invalid_level_indexes() {
let dir = TempDir::new().unwrap();
let sst_dir = dir.path().join("sst");
std::fs::create_dir_all(&sst_dir).unwrap();
let file = make_live_sst(&sst_dir, 1, b"a", b"z");
let mut vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
let kind = match vs.apply(&[VersionEdit::AddFile {
level: MAX_LEVELS,
file,
}]) {
Err(e) => e.kind(),
Ok(_) => panic!("expected invalid level error"),
};
assert_eq!(kind, io::ErrorKind::InvalidData);
assert_eq!(vs.current().levels.iter().map(Vec::len).sum::<usize>(), 0);
let kind = match vs.apply(&[VersionEdit::RemoveFile {
level: MAX_LEVELS,
file_id: 1,
}]) {
Err(e) => e.kind(),
Ok(_) => panic!("expected invalid level error"),
};
assert_eq!(kind, io::ErrorKind::InvalidData);
}
#[test]
fn open_rejects_manifest_with_invalid_level_index() {
let dir = TempDir::new().unwrap();
let sst_dir = dir.path().join("sst");
std::fs::create_dir_all(&sst_dir).unwrap();
let records = [ManifestRecord::AddFile {
level: MAX_LEVELS,
meta: test_meta(1),
}];
std::fs::write(
dir.path().join("MANIFEST"),
VersionSet::encode_records(&records),
)
.unwrap();
let kind = match VersionSet::open(dir.path(), &sst_dir) {
Err(e) => e.kind(),
Ok(_) => panic!("expected invalid level error"),
};
assert_eq!(kind, io::ErrorKind::InvalidData);
}
#[test]
fn replay_survives_truncated_trailer() {
let dir = TempDir::new().unwrap();
let sst_dir = dir.path().join("sst");
std::fs::create_dir_all(&sst_dir).unwrap();
let file1 = make_live_sst(&sst_dir, 1, b"a", b"m");
let mut vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
vs.apply(&[VersionEdit::AddFile {
level: 0,
file: Arc::clone(&file1),
}])
.unwrap();
vs.apply(&[VersionEdit::SetLastSeq(50)]).unwrap();
drop(vs);
let path = dir.path().join("MANIFEST");
let current = std::fs::metadata(&path).unwrap().len();
OpenOptions::new()
.write(true)
.open(&path)
.unwrap()
.set_len(current - 2)
.unwrap();
let vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
let v = vs.current();
assert!(v.levels[0].len() <= 1);
}
#[test]
fn reset_record_clears_files_and_sets_wal_floor_atomically() {
let dir = TempDir::new().unwrap();
let sst_dir = dir.path().join("sst");
std::fs::create_dir_all(&sst_dir).unwrap();
let file1 = make_live_sst(&sst_dir, 1, b"a", b"m");
let file2 = make_live_sst(&sst_dir, 2, b"n", b"z");
{
let mut vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
vs.apply(&[
VersionEdit::AddFile {
level: 0,
file: Arc::clone(&file1),
},
VersionEdit::AddFile {
level: 1,
file: Arc::clone(&file2),
},
VersionEdit::SetLastSeq(50),
VersionEdit::SetNextFileId(9),
])
.unwrap();
vs.apply(&[VersionEdit::Reset {
next_file_id: 10,
min_wal_id: 9,
}])
.unwrap();
}
let vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
let v = vs.current();
assert!(v.levels.iter().all(Vec::is_empty));
assert_eq!(v.last_seq, 0);
assert_eq!(v.next_file_id, 10);
assert_eq!(v.min_wal_id, 9);
}
#[test]
fn open_truncates_truncated_manifest_tail_before_append() {
let dir = TempDir::new().unwrap();
let sst_dir = dir.path().join("sst");
std::fs::create_dir_all(&sst_dir).unwrap();
{
let mut vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
vs.apply(&[VersionEdit::SetLastSeq(7)]).unwrap();
vs.apply(&[VersionEdit::SetLastSeq(11)]).unwrap();
}
let path = dir.path().join("MANIFEST");
let current = std::fs::metadata(&path).unwrap().len();
OpenOptions::new()
.write(true)
.open(&path)
.unwrap()
.set_len(current - 2)
.unwrap();
{
let mut vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
assert_eq!(vs.current().last_seq, 7);
vs.apply(&[VersionEdit::SetLastSeq(99)]).unwrap();
}
let vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
assert_eq!(vs.current().last_seq, 99);
}
#[test]
fn open_truncates_corrupt_manifest_tail_before_append() {
let dir = TempDir::new().unwrap();
let sst_dir = dir.path().join("sst");
std::fs::create_dir_all(&sst_dir).unwrap();
{
let mut vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
vs.apply(&[VersionEdit::SetLastSeq(7)]).unwrap();
vs.apply(&[VersionEdit::SetLastSeq(11)]).unwrap();
}
let path = dir.path().join("MANIFEST");
let checksum_offset = second_record_checksum_offset(&path);
let mut data = std::fs::read(&path).unwrap();
data[checksum_offset] ^= 0xFF;
std::fs::write(&path, data).unwrap();
{
let mut vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
assert_eq!(vs.current().last_seq, 7);
vs.apply(&[VersionEdit::SetLastSeq(99)]).unwrap();
}
let vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
assert_eq!(vs.current().last_seq, 99);
}
#[test]
fn a_long_running_manifest_stays_bounded_instead_of_growing_forever() {
let dir = TempDir::new().unwrap();
let sst_dir = dir.path().join("sst");
std::fs::create_dir_all(&sst_dir).unwrap();
let path = dir.path().join("MANIFEST");
let mut vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
const EDITS: u64 = 8_000;
for id in 1..=EDITS {
vs.apply(&[VersionEdit::SetNextFileId(id)]).unwrap();
}
drop(vs);
let len = std::fs::metadata(&path).unwrap().len();
let bound = VersionSet::compact_threshold(0);
assert!(
len <= bound,
"manifest grew to {len} bytes against a {bound}-byte bound: \
an append-only log that is never rewritten grows without limit"
);
let reopened = VersionSet::open(dir.path(), &sst_dir).unwrap();
assert_eq!(reopened.current().next_file_id, EDITS);
}
#[test]
fn a_manifest_from_a_newer_format_is_refused() {
let mut stamp = VersionSet::encode_stamp();
stamp[7] = MANIFEST_FORMAT_V1 + 1;
let checksum = checksum::manifest_record(0, &stamp[0..8]);
stamp[8..12].copy_from_slice(&checksum.to_le_bytes());
let err = VersionSet::stamp_len(&stamp).expect_err("a newer format must not be parsed");
assert!(err.to_string().contains("newer regolith"), "{err}");
}
#[test]
fn compact_manifest_rewrites_to_canonical_form() {
let dir = TempDir::new().unwrap();
let sst_dir = dir.path().join("sst");
std::fs::create_dir_all(&sst_dir).unwrap();
let file1 = make_live_sst(&sst_dir, 1, b"a", b"c");
let file2 = make_live_sst(&sst_dir, 2, b"d", b"f");
let mut vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
vs.apply(&[VersionEdit::Reset {
next_file_id: 7,
min_wal_id: 7,
}])
.unwrap();
vs.apply(&[VersionEdit::AddFile {
level: 0,
file: Arc::clone(&file1),
}])
.unwrap();
vs.apply(&[VersionEdit::AddFile {
level: 1,
file: Arc::clone(&file2),
}])
.unwrap();
vs.apply(&[VersionEdit::SetLastSeq(500), VersionEdit::SetNextFileId(99)])
.unwrap();
let pre_size = std::fs::metadata(dir.path().join("MANIFEST"))
.unwrap()
.len();
vs.compact_manifest().unwrap();
let post_size = std::fs::metadata(dir.path().join("MANIFEST"))
.unwrap()
.len();
assert!(post_size <= pre_size + 64);
drop(vs);
let vs = VersionSet::open(dir.path(), &sst_dir).unwrap();
let v = vs.current();
assert_eq!(v.levels[0].len(), 1);
assert_eq!(v.levels[1].len(), 1);
assert_eq!(v.last_seq, 500);
assert_eq!(v.next_file_id, 99);
assert_eq!(v.min_wal_id, 7);
}
}