use async_trait::async_trait;
use corium_core::{
Datom, EntityId,
encoding::{decode_value, encode_value},
};
use std::{
collections::HashMap,
fs::{self, File, OpenOptions},
io::{self, Write},
path::{Path, PathBuf},
sync::{Arc, Mutex, RwLock},
};
use thiserror::Error;
const CHECKSUMMED_FRAME: u64 = 1 << 63;
const FRAME_CHECKSUM_LEN: usize = size_of::<u32>();
const RANGE_READ_CHUNK_BYTES: u64 = 4 * 1024 * 1024;
const MAX_CACHED_READ_VERSION_FILES: usize = 8;
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct TxRecord {
pub t: u64,
pub tx_instant: i64,
pub datoms: Vec<Datom>,
}
#[derive(Debug, Error)]
pub enum LogError {
#[error("log I/O failed: {0}")]
Io(#[from] io::Error),
#[error("corrupt transaction log")]
Corrupt,
#[error("native transaction log store failed: {0}")]
Native(String),
#[error("this transaction log requires asynchronous access")]
AsyncOnly,
}
#[async_trait]
pub trait TransactionLog: Send + Sync {
fn append(&self, record: &TxRecord) -> Result<(), LogError>;
async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
self.append(record)
}
async fn append_batch_async(&self, records: &[TxRecord]) -> Result<(), LogError> {
for record in records {
self.append_async(record).await?;
}
Ok(())
}
fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError>;
async fn tx_range_async(
&self,
start: u64,
end: Option<u64>,
) -> Result<Vec<TxRecord>, LogError> {
self.tx_range(start, end)
}
fn replay(&self) -> Result<Vec<TxRecord>, LogError> {
self.tx_range(0, None)
}
async fn replay_async(&self) -> Result<Vec<TxRecord>, LogError> {
self.tx_range_async(0, None).await
}
}
#[derive(Clone, Default)]
pub struct MemoryLog(Arc<RwLock<Vec<TxRecord>>>);
impl TransactionLog for MemoryLog {
fn append(&self, record: &TxRecord) -> Result<(), LogError> {
let mut records = self.0.write().expect("poisoned log lock");
if records.last().map_or(1, |r| r.t + 1) != record.t {
return Err(LogError::Corrupt);
}
records.push(record.clone());
Ok(())
}
fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
Ok(self
.0
.read()
.expect("poisoned log lock")
.iter()
.filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
.cloned()
.collect())
}
}
pub struct FileLog {
state: RwLock<IndexedFile>,
}
impl FileLog {
pub fn open(path: impl AsRef<Path>) -> Result<Self, LogError> {
let path = path.as_ref().to_path_buf();
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}
let file = IndexedFile::open(&path, true, true)?;
file.validate_contiguous_prefix(file.frames.len())?;
Ok(Self {
state: RwLock::new(file),
})
}
}
impl TransactionLog for FileLog {
fn append(&self, record: &TxRecord) -> Result<(), LogError> {
let mut state = self.state.write().expect("poisoned log lock");
state.refresh()?;
state.validate_contiguous_prefix(state.frames.len())?;
if next_t(&state.frames)? != record.t {
return Err(LogError::Corrupt);
}
state.append(record)
}
fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
if end.is_some_and(|end| end <= start) {
return Ok(Vec::new());
}
let indexed = {
let state = self.state.read().expect("poisoned log lock");
range_is_indexed(&state.frames, end)?
};
if !indexed {
let mut state = self.state.write().expect("poisoned log lock");
state.refresh()?;
state.validate_contiguous_prefix(state.frames.len())?;
}
self.state
.read()
.expect("poisoned log lock")
.tx_range(start, end)
}
}
pub struct VersionedLog {
dir: PathBuf,
name: String,
state: RwLock<VersionedLogState>,
}
struct VersionedLogState {
files: Vec<VersionedFile>,
write_version: Option<u64>,
next_t: u64,
}
struct VersionedFile {
version: u64,
file: IndexedFile,
}
impl VersionedLog {
pub fn open(dir: impl AsRef<Path>, name: &str, write_version: u64) -> Result<Self, LogError> {
let dir = dir.as_ref().to_path_buf();
fs::create_dir_all(&dir)?;
let write_path = version_path(&dir, name, write_version);
let mut files = Vec::new();
for (version, path) in version_files(&dir, name) {
let writable = version == write_version;
files.push(VersionedFile {
version,
file: IndexedFile::open(&path, writable, writable)?,
});
close_cold_version_files(&mut files);
}
if !files.iter().any(|file| file.version == write_version) {
files.push(VersionedFile {
version: write_version,
file: IndexedFile::open(&write_path, true, true)?,
});
files.sort_by_key(|file| file.version);
}
close_cold_version_files(&mut files);
let cutoffs = validated_version_cutoffs(&files)?;
let next_t = merged_next_t(&files, &cutoffs)?;
Ok(Self {
dir,
name: name.to_owned(),
state: RwLock::new(VersionedLogState {
files,
write_version: Some(write_version),
next_t,
}),
})
}
pub fn open_read_only(dir: impl AsRef<Path>, name: &str) -> Result<Self, LogError> {
let dir = dir.as_ref().to_path_buf();
let mut files = open_version_files(&dir, name)?;
close_cold_version_files(&mut files);
let cutoffs = validated_version_cutoffs(&files)?;
Ok(Self {
name: name.to_owned(),
state: RwLock::new(VersionedLogState {
next_t: merged_next_t(&files, &cutoffs)?,
write_version: None,
files,
}),
dir,
})
}
#[must_use]
pub fn exists(dir: impl AsRef<Path>, name: &str) -> bool {
!version_files(dir.as_ref(), name).is_empty()
}
pub fn delete_all(dir: impl AsRef<Path>, name: &str) -> Result<(), LogError> {
for (_, path) in version_files(dir.as_ref(), name) {
match fs::remove_file(&path) {
Ok(()) => {}
Err(error) if error.kind() == io::ErrorKind::NotFound => {}
Err(error) => return Err(error.into()),
}
}
Ok(())
}
}
impl TransactionLog for VersionedLog {
fn append(&self, record: &TxRecord) -> Result<(), LogError> {
let mut state = self.state.write().expect("poisoned log lock");
if state.next_t != record.t {
return Err(LogError::Corrupt);
}
let write_version = state
.write_version
.ok_or_else(|| LogError::Native("transaction log is read-only".into()))?;
let write_index = state
.files
.iter()
.position(|file| file.version == write_version)
.ok_or(LogError::Corrupt)?;
let cutoffs = version_cutoffs(&state.files);
if cutoffs[write_index] == u64::MAX
&& !state.files[write_index].file.frames.is_empty()
&& next_t(&state.files[write_index].file.frames)? != record.t
{
return Err(LogError::Corrupt);
}
state.files[write_index].file.append(record)?;
state.next_t = state.next_t.checked_add(1).ok_or(LogError::Corrupt)?;
Ok(())
}
fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
if end.is_some_and(|end| end <= start) {
return Ok(Vec::new());
}
{
let state = self.state.read().expect("poisoned log lock");
let cutoffs = validated_version_cutoffs(&state.files)?;
if range_is_merged_indexed(&state.files, &cutoffs, end)? {
return read_merged_range(&state.files, &cutoffs, start, end);
}
}
{
let mut state = self.state.write().expect("poisoned log lock");
refresh_version_files(&self.dir, &self.name, &mut state.files)?;
close_cold_version_files(&mut state.files);
validated_version_cutoffs(&state.files)?;
}
let state = self.state.read().expect("poisoned log lock");
let cutoffs = validated_version_cutoffs(&state.files)?;
read_merged_range(&state.files, &cutoffs, start, end)
}
}
fn merge_versions(mut per_version: Vec<Vec<TxRecord>>) -> Vec<TxRecord> {
let mut cutoff = u64::MAX;
for records in per_version.iter_mut().rev() {
let first = records.first().map(|r| r.t);
records.retain(|r| r.t < cutoff);
if let Some(first) = first {
cutoff = cutoff.min(first);
}
}
per_version.into_iter().flatten().collect()
}
#[async_trait]
pub trait NativeLogStorage: Send + Sync {
async fn put_batch(
&self,
name: &str,
version: u64,
records: &[(u64, Vec<u8>)],
) -> Result<bool, LogError>;
async fn read_record(
&self,
name: &str,
version: u64,
t: u64,
) -> Result<Option<Vec<u8>>, LogError>;
async fn list_records(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
async fn read_legacy_chunk(
&self,
name: &str,
version: u64,
chunk: u64,
) -> Result<Option<Vec<u8>>, LogError>;
async fn list_legacy_chunks(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
async fn delete_all(&self, name: &str) -> Result<(), LogError>;
}
pub struct NativeVersionedLog<S: ?Sized> {
storage: Arc<S>,
name: String,
write_version: u64,
read_only: bool,
next_t: tokio::sync::Mutex<u64>,
}
impl<S: NativeLogStorage + ?Sized + 'static> NativeVersionedLog<S> {
pub async fn open(storage: Arc<S>, name: &str, write_version: u64) -> Result<Self, LogError> {
let records = read_native_merged(storage.as_ref(), name).await?;
let next_t = records.last().map_or(1, |r| r.t + 1);
Ok(Self {
storage,
name: name.to_owned(),
write_version,
read_only: false,
next_t: tokio::sync::Mutex::new(next_t),
})
}
#[must_use]
pub fn open_read_only(storage: Arc<S>, name: &str) -> Self {
Self {
storage,
name: name.to_owned(),
write_version: 0,
read_only: true,
next_t: tokio::sync::Mutex::new(0),
}
}
}
#[async_trait]
impl<S: NativeLogStorage + ?Sized + 'static> TransactionLog for NativeVersionedLog<S> {
fn append(&self, record: &TxRecord) -> Result<(), LogError> {
let _ = record;
Err(LogError::AsyncOnly)
}
async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
self.append_batch_async(std::slice::from_ref(record)).await
}
async fn append_batch_async(&self, records: &[TxRecord]) -> Result<(), LogError> {
if self.read_only {
return Err(LogError::Native("transaction log is read-only".into()));
}
if records.is_empty() {
return Ok(());
}
let mut next_t = self.next_t.lock().await;
for (offset, record) in records.iter().enumerate() {
if record.t != *next_t + offset as u64 {
return Err(LogError::Corrupt);
}
}
let framed = records
.iter()
.map(|record| {
let mut bytes = Vec::new();
append_framed_record(&mut bytes, record)?;
Ok((record.t, bytes))
})
.collect::<Result<Vec<_>, LogError>>()?;
if !self
.storage
.put_batch(&self.name, self.write_version, &framed)
.await?
{
return Err(LogError::Corrupt);
}
*next_t += records.len() as u64;
Ok(())
}
fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
let _ = (start, end);
Err(LogError::AsyncOnly)
}
async fn tx_range_async(
&self,
start: u64,
end: Option<u64>,
) -> Result<Vec<TxRecord>, LogError> {
let _guard = self.next_t.lock().await;
Ok(read_native_merged(self.storage.as_ref(), &self.name)
.await?
.into_iter()
.filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
.collect())
}
}
async fn read_native_merged<S: NativeLogStorage + ?Sized>(
storage: &S,
name: &str,
) -> Result<Vec<TxRecord>, LogError> {
use std::collections::BTreeMap;
let mut per_version: BTreeMap<u64, Vec<TxRecord>> = BTreeMap::new();
let mut chunks = storage.list_legacy_chunks(name).await?;
chunks.sort_unstable();
for (version, chunk) in chunks {
let bytes = storage
.read_legacy_chunk(name, version, chunk)
.await?
.unwrap_or_default();
per_version
.entry(version)
.or_default()
.extend(decode_framed_records(&bytes)?);
}
let mut records = storage.list_records(name).await?;
records.sort_unstable();
for (version, t) in records {
let bytes = storage
.read_record(name, version, t)
.await?
.unwrap_or_default();
per_version
.entry(version)
.or_default()
.extend(decode_framed_records(&bytes)?);
}
let per_version: Vec<Vec<TxRecord>> = per_version
.into_values()
.map(|mut records| {
records.sort_by_key(|record| record.t);
records
})
.collect();
let merged = merge_versions(per_version);
for pair in merged.windows(2) {
if pair[1].t != pair[0].t + 1 {
return Err(LogError::Corrupt);
}
}
Ok(merged)
}
type VersionedRecords = Arc<Mutex<Vec<(u64, TxRecord)>>>;
#[derive(Clone, Default)]
pub struct MemLogRegistry {
logs: Arc<Mutex<HashMap<String, VersionedRecords>>>,
}
impl MemLogRegistry {
#[must_use]
pub fn new() -> Self {
Self::default()
}
fn entry(&self, name: &str) -> VersionedRecords {
Arc::clone(
self.logs
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.entry(name.to_owned())
.or_default(),
)
}
#[must_use]
pub fn open(&self, name: &str, write_version: u64) -> MemVersionedLog {
let records = self.entry(name);
let next_t = {
let guard = records
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
MemVersionedLog::merged(&guard)
.last()
.map_or(1, |r| r.t + 1)
};
MemVersionedLog {
records,
write_version,
next_t: Mutex::new(next_t),
}
}
#[must_use]
pub fn exists(&self, name: &str) -> bool {
self.logs
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(name)
.is_some_and(|entry| {
!entry
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_empty()
})
}
pub fn delete_all(&self, name: &str) {
self.logs
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(name);
}
}
pub struct MemVersionedLog {
records: VersionedRecords,
write_version: u64,
next_t: Mutex<u64>,
}
impl MemVersionedLog {
fn merged(records: &[(u64, TxRecord)]) -> Vec<TxRecord> {
let mut versions: Vec<u64> = records.iter().map(|(version, _)| *version).collect();
versions.sort_unstable();
versions.dedup();
let per_version = versions
.into_iter()
.map(|version| {
records
.iter()
.filter(|(record_version, _)| *record_version == version)
.map(|(_, record)| record.clone())
.collect::<Vec<_>>()
})
.collect();
merge_versions(per_version)
}
}
impl TransactionLog for MemVersionedLog {
fn append(&self, record: &TxRecord) -> Result<(), LogError> {
let mut next_t = self
.next_t
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if *next_t != record.t {
return Err(LogError::Corrupt);
}
self.records
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push((self.write_version, record.clone()));
*next_t += 1;
Ok(())
}
fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
let records = self
.records
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
Ok(Self::merged(&records)
.into_iter()
.filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
.collect())
}
}
#[derive(Clone, Copy)]
struct FrameIndex {
t: u64,
offset: u64,
len: u64,
}
struct IndexedFile {
path: PathBuf,
file: Option<Arc<File>>,
writable: bool,
poisoned: bool,
frames: Vec<FrameIndex>,
durable_len: u64,
first_gap: Option<usize>,
}
impl IndexedFile {
fn open(path: &Path, writable: bool, truncate_torn: bool) -> Result<Self, LogError> {
let file = Arc::new(open_index_file(path, writable)?);
let (frames, durable_len) = scan_frames(file.as_ref(), 0)?;
validate_sorted_frames(&frames)?;
if truncate_torn && file.metadata()?.len() > durable_len {
file.set_len(durable_len)?;
file.sync_all()?;
}
Ok(Self {
path: path.to_path_buf(),
file: Some(file),
writable,
poisoned: false,
first_gap: first_gap_index(&frames),
frames,
durable_len,
})
}
fn refresh(&mut self) -> Result<(), LogError> {
self.ensure_healthy()?;
let file_len = self.physical_len()?;
if file_len < self.durable_len {
let file = self.open_for_read()?;
let (frames, durable_len) = scan_frames(file.as_ref(), 0)?;
validate_sorted_frames(&frames)?;
self.first_gap = first_gap_index(&frames);
self.frames = frames;
self.durable_len = durable_len;
} else if file_len > self.durable_len {
let file = self.open_for_read()?;
let (new_frames, durable_len) = scan_frames(file.as_ref(), self.durable_len)?;
validate_sorted_extension(&self.frames, &new_frames)?;
let existing_len = self.frames.len();
if self.first_gap.is_none() {
self.first_gap =
extension_first_gap(&self.frames, &new_frames).map(|gap| existing_len + gap);
}
self.frames.extend(new_frames);
self.durable_len = durable_len;
}
Ok(())
}
fn append(&mut self, record: &TxRecord) -> Result<(), LogError> {
self.ensure_healthy()?;
if !self.writable {
return Err(LogError::Native("transaction log is read-only".into()));
}
let file = self.open_for_read()?;
if file.metadata()?.len() != self.durable_len {
return Err(LogError::Corrupt);
}
let mut frame = Vec::new();
append_framed_record(&mut frame, record)?;
let frame_len = u64::try_from(frame.len()).map_err(|_| LogError::Corrupt)?;
let offset = self.durable_len;
let mut writer = file.as_ref();
if let Err(error) = writer.write_all(&frame) {
self.poisoned = true;
return Err(error.into());
}
if let Err(error) = file.sync_all() {
self.poisoned = true;
return Err(error.into());
}
self.durable_len = self
.durable_len
.checked_add(frame_len)
.ok_or(LogError::Corrupt)?;
if self.first_gap.is_none()
&& self
.frames
.last()
.is_some_and(|previous| previous.t.checked_add(1) != Some(record.t))
{
self.first_gap = Some(self.frames.len());
}
self.frames.push(FrameIndex {
t: record.t,
offset,
len: frame_len,
});
Ok(())
}
fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
self.ensure_healthy()?;
let first = self.frames.partition_point(|frame| frame.t < start);
let last = end.map_or(self.frames.len(), |end| {
self.frames.partition_point(|frame| frame.t < end)
});
if first >= last {
return Ok(Vec::new());
}
let file = self.open_for_read()?;
let mut records = Vec::with_capacity(last - first);
let mut chunk_first = first;
while chunk_first < last {
let offset = self.frames[chunk_first].offset;
let mut chunk_last = chunk_first + 1;
while chunk_last < last {
let candidate_end = frame_end(self.frames[chunk_last])?;
if candidate_end.checked_sub(offset).ok_or(LogError::Corrupt)?
> RANGE_READ_CHUNK_BYTES
{
break;
}
chunk_last += 1;
}
let byte_end = frame_end(self.frames[chunk_last - 1])?;
let byte_len = usize::try_from(byte_end.checked_sub(offset).ok_or(LogError::Corrupt)?)
.map_err(|_| LogError::Corrupt)?;
let mut bytes = vec![0; byte_len];
read_exact_at(file.as_ref(), &mut bytes, offset)?;
let chunk_records = decode_framed_records(&bytes)?;
if chunk_records.len() != chunk_last - chunk_first
|| chunk_records
.iter()
.zip(&self.frames[chunk_first..chunk_last])
.any(|(record, frame)| record.t != frame.t)
{
return Err(LogError::Corrupt);
}
records.extend(chunk_records);
chunk_first = chunk_last;
}
Ok(records)
}
fn validate_contiguous_prefix(&self, retained: usize) -> Result<(), LogError> {
if self.first_gap.is_some_and(|gap| gap < retained) {
return Err(LogError::Corrupt);
}
Ok(())
}
fn ensure_healthy(&self) -> Result<(), LogError> {
if self.poisoned {
return Err(io::Error::other("transaction log handle is poisoned").into());
}
Ok(())
}
fn open_for_read(&self) -> Result<Arc<File>, LogError> {
self.file.as_ref().map_or_else(
|| Ok(Arc::new(open_index_file(&self.path, false)?)),
|file| Ok(Arc::clone(file)),
)
}
fn physical_len(&self) -> Result<u64, LogError> {
Ok(self
.file
.as_ref()
.map_or_else(|| fs::metadata(&self.path), |file| file.metadata())?
.len())
}
fn close_cached_reader(&mut self) {
if !self.writable {
self.file = None;
}
}
}
fn open_index_file(path: &Path, writable: bool) -> Result<File, io::Error> {
let mut options = OpenOptions::new();
options.read(true);
if writable {
options.create(true).write(true).append(true);
}
options.open(path)
}
fn validate_sorted_frames(frames: &[FrameIndex]) -> Result<(), LogError> {
for pair in frames.windows(2) {
if pair[0].t >= pair[1].t {
return Err(LogError::Corrupt);
}
}
Ok(())
}
fn validate_sorted_extension(
existing: &[FrameIndex],
appended: &[FrameIndex],
) -> Result<(), LogError> {
validate_sorted_frames(appended)?;
if let (Some(previous), Some(next)) = (existing.last(), appended.first())
&& previous.t >= next.t
{
return Err(LogError::Corrupt);
}
Ok(())
}
fn first_gap_index(frames: &[FrameIndex]) -> Option<usize> {
frames
.windows(2)
.position(|pair| pair[0].t.checked_add(1) != Some(pair[1].t))
.map(|index| index + 1)
}
fn extension_first_gap(existing: &[FrameIndex], appended: &[FrameIndex]) -> Option<usize> {
if let (Some(previous), Some(next)) = (existing.last(), appended.first())
&& previous.t.checked_add(1) != Some(next.t)
{
return Some(0);
}
first_gap_index(appended)
}
fn frame_end(frame: FrameIndex) -> Result<u64, LogError> {
frame.offset.checked_add(frame.len).ok_or(LogError::Corrupt)
}
fn next_t(frames: &[FrameIndex]) -> Result<u64, LogError> {
frames.last().map_or(Ok(1), |frame| {
frame.t.checked_add(1).ok_or(LogError::Corrupt)
})
}
fn range_is_indexed(frames: &[FrameIndex], end: Option<u64>) -> Result<bool, LogError> {
end.map_or(Ok(false), |end| Ok(end <= next_t(frames)?))
}
fn version_path(dir: &Path, name: &str, version: u64) -> PathBuf {
if version == 0 {
dir.join(format!("{name}.log"))
} else {
dir.join(format!("{name}.v{version}.log"))
}
}
fn version_files(dir: &Path, name: &str) -> Vec<(u64, PathBuf)> {
let mut files = Vec::new();
let legacy = version_path(dir, name, 0);
if legacy.is_file() {
files.push((0, legacy));
}
let prefix = format!("{name}.v");
if let Ok(entries) = fs::read_dir(dir) {
for entry in entries.flatten() {
let file_name = entry.file_name();
let Some(text) = file_name.to_str() else {
continue;
};
if let Some(version) = text
.strip_prefix(&prefix)
.and_then(|rest| rest.strip_suffix(".log"))
.and_then(|v| v.parse::<u64>().ok())
&& version > 0
{
files.push((version, entry.path()));
}
}
}
files.sort_by_key(|(version, _)| *version);
files
}
fn open_version_files(dir: &Path, name: &str) -> Result<Vec<VersionedFile>, LogError> {
let mut files = Vec::new();
for (version, path) in version_files(dir, name) {
files.push(VersionedFile {
version,
file: IndexedFile::open(&path, false, false)?,
});
close_cold_version_files(&mut files);
}
Ok(files)
}
fn refresh_version_files(
dir: &Path,
name: &str,
files: &mut Vec<VersionedFile>,
) -> Result<(), LogError> {
for file in &mut *files {
file.file.refresh()?;
}
for (version, path) in version_files(dir, name) {
if files.iter().all(|file| file.version != version) {
files.push(VersionedFile {
version,
file: IndexedFile::open(&path, false, false)?,
});
close_cold_version_files(files);
}
}
files.sort_by_key(|file| file.version);
Ok(())
}
fn close_cold_version_files(files: &mut [VersionedFile]) {
let mut cached_readers = 0;
for file in files.iter_mut().rev() {
if file.file.writable {
continue;
}
if file.file.file.is_some() {
if cached_readers < MAX_CACHED_READ_VERSION_FILES {
cached_readers += 1;
} else {
file.file.close_cached_reader();
}
}
}
}
fn version_cutoffs(files: &[VersionedFile]) -> Vec<u64> {
let mut cutoffs = vec![u64::MAX; files.len()];
let mut cutoff = u64::MAX;
for (index, file) in files.iter().enumerate().rev() {
cutoffs[index] = cutoff;
if let Some(first) = file.file.frames.first() {
cutoff = cutoff.min(first.t);
}
}
cutoffs
}
fn validated_version_cutoffs(files: &[VersionedFile]) -> Result<Vec<u64>, LogError> {
let cutoffs = version_cutoffs(files);
let mut previous_t: Option<u64> = None;
for (file, cutoff) in files.iter().zip(&cutoffs) {
let retained = file.file.frames.partition_point(|frame| frame.t < *cutoff);
if retained == 0 {
continue;
}
file.file.validate_contiguous_prefix(retained)?;
let first_t = file.file.frames[0].t;
if previous_t.is_some_and(|previous| previous.checked_add(1) != Some(first_t)) {
return Err(LogError::Corrupt);
}
previous_t = Some(file.file.frames[retained - 1].t);
}
Ok(cutoffs)
}
fn merged_next_t(files: &[VersionedFile], cutoffs: &[u64]) -> Result<u64, LogError> {
for (file, cutoff) in files.iter().zip(cutoffs).rev() {
let retained = file.file.frames.partition_point(|frame| frame.t < *cutoff);
if retained > 0 {
return file.file.frames[retained - 1]
.t
.checked_add(1)
.ok_or(LogError::Corrupt);
}
}
Ok(1)
}
fn range_is_merged_indexed(
files: &[VersionedFile],
cutoffs: &[u64],
end: Option<u64>,
) -> Result<bool, LogError> {
end.map_or(Ok(false), |end| Ok(end <= merged_next_t(files, cutoffs)?))
}
fn read_merged_range(
files: &[VersionedFile],
cutoffs: &[u64],
start: u64,
end: Option<u64>,
) -> Result<Vec<TxRecord>, LogError> {
let mut records = Vec::new();
for (file, cutoff) in files.iter().zip(cutoffs) {
let end = Some(end.map_or(*cutoff, |end| end.min(*cutoff)));
records.extend(file.file.tx_range(start, end)?);
}
Ok(records)
}
fn encode_record(record: &TxRecord) -> Vec<u8> {
let mut out = Vec::new();
out.extend_from_slice(&record.t.to_be_bytes());
out.extend_from_slice(&record.tx_instant.to_be_bytes());
out.extend_from_slice(&(record.datoms.len() as u64).to_be_bytes());
for d in &record.datoms {
out.extend_from_slice(&d.e.raw().to_be_bytes());
out.extend_from_slice(&d.a.raw().to_be_bytes());
out.extend_from_slice(&d.tx.raw().to_be_bytes());
out.push(u8::from(d.added));
let v = encode_value(&d.v);
out.extend_from_slice(&(v.len() as u64).to_be_bytes());
out.extend_from_slice(&v);
}
out
}
fn decode_record(mut bytes: &[u8]) -> Result<TxRecord, LogError> {
fn take<'a>(bytes: &mut &'a [u8], n: usize) -> Result<&'a [u8], LogError> {
let value = bytes.get(..n).ok_or(LogError::Corrupt)?;
*bytes = &bytes[n..];
Ok(value)
}
fn u64_be(bytes: &mut &[u8]) -> Result<u64, LogError> {
Ok(u64::from_be_bytes(
take(bytes, 8)?.try_into().map_err(|_| LogError::Corrupt)?,
))
}
let t = u64_be(&mut bytes)?;
let tx_instant = i64::from_be_bytes(
take(&mut bytes, 8)?
.try_into()
.map_err(|_| LogError::Corrupt)?,
);
let count = u64_be(&mut bytes)?;
let mut datoms = Vec::new();
for _ in 0..count {
let e = EntityId::from_raw(u64_be(&mut bytes)?);
let a = EntityId::from_raw(u64_be(&mut bytes)?);
let tx = EntityId::from_raw(u64_be(&mut bytes)?);
let added = take(&mut bytes, 1)?[0] != 0;
let len = usize::try_from(u64_be(&mut bytes)?).map_err(|_| LogError::Corrupt)?;
let raw = take(&mut bytes, len)?;
let (v, used) = decode_value(raw).map_err(|_| LogError::Corrupt)?;
if used != len {
return Err(LogError::Corrupt);
}
datoms.push(Datom { e, a, v, tx, added });
}
if !bytes.is_empty() {
return Err(LogError::Corrupt);
}
Ok(TxRecord {
t,
tx_instant,
datoms,
})
}
fn frame_header(payload_len: usize) -> Result<[u8; 8], LogError> {
let payload_len = u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?;
if payload_len & CHECKSUMMED_FRAME != 0 {
return Err(LogError::Corrupt);
}
Ok((payload_len | CHECKSUMMED_FRAME).to_be_bytes())
}
fn frame_payload_len(header: [u8; 8]) -> Result<(usize, bool), LogError> {
let encoded = u64::from_be_bytes(header);
let checksummed = encoded & CHECKSUMMED_FRAME != 0;
let payload_len = encoded & !CHECKSUMMED_FRAME;
Ok((
usize::try_from(payload_len).map_err(|_| LogError::Corrupt)?,
checksummed,
))
}
fn frame_checksum(header: [u8; 8], payload: &[u8]) -> u32 {
crc32c::crc32c_append(crc32c::crc32c(&header), payload)
}
#[cfg(unix)]
fn read_exact_at(file: &File, mut bytes: &mut [u8], mut offset: u64) -> io::Result<()> {
use std::os::unix::fs::FileExt;
while !bytes.is_empty() {
match file.read_at(bytes, offset) {
Ok(0) => return Err(io::ErrorKind::UnexpectedEof.into()),
Ok(read) => {
offset = offset
.checked_add(u64::try_from(read).expect("read length fits u64"))
.ok_or_else(|| io::Error::other("file offset overflow"))?;
bytes = &mut bytes[read..];
}
Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
Err(error) => return Err(error),
}
}
Ok(())
}
#[cfg(windows)]
fn read_exact_at(file: &File, mut bytes: &mut [u8], mut offset: u64) -> io::Result<()> {
use std::os::windows::fs::FileExt;
while !bytes.is_empty() {
match file.seek_read(bytes, offset) {
Ok(0) => return Err(io::ErrorKind::UnexpectedEof.into()),
Ok(read) => {
offset = offset
.checked_add(u64::try_from(read).expect("read length fits u64"))
.ok_or_else(|| io::Error::other("file offset overflow"))?;
bytes = &mut bytes[read..];
}
Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
Err(error) => return Err(error),
}
}
Ok(())
}
#[cfg(not(any(unix, windows)))]
fn read_exact_at(file: &File, bytes: &mut [u8], offset: u64) -> io::Result<()> {
use std::io::{Read, Seek, SeekFrom};
let mut file = file.try_clone()?;
file.seek(SeekFrom::Start(offset))?;
file.read_exact(bytes)
}
fn scan_frames(file: &File, offset: u64) -> Result<(Vec<FrameIndex>, u64), LogError> {
let file_len = file.metadata()?.len();
let mut frames = Vec::new();
let mut durable_len = offset;
loop {
if file_len.saturating_sub(durable_len) < 8 {
break;
}
let mut len = [0; 8];
read_exact_at(file, &mut len, durable_len)?;
let (payload_len, checksummed) = frame_payload_len(len)?;
let frame_len = 8_u64
.checked_add(u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?)
.and_then(|len| {
len.checked_add(if checksummed {
u64::try_from(FRAME_CHECKSUM_LEN).expect("checksum length fits u64")
} else {
0
})
})
.ok_or(LogError::Corrupt)?;
if file_len.saturating_sub(durable_len) < frame_len {
break;
}
let mut payload = vec![0; payload_len];
let payload_offset = durable_len.checked_add(8).ok_or(LogError::Corrupt)?;
read_exact_at(file, &mut payload, payload_offset)?;
if checksummed {
let mut stored_checksum = [0; FRAME_CHECKSUM_LEN];
let checksum_offset = payload_offset
.checked_add(u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?)
.ok_or(LogError::Corrupt)?;
read_exact_at(file, &mut stored_checksum, checksum_offset)?;
if u32::from_be_bytes(stored_checksum) != frame_checksum(len, &payload) {
return Err(LogError::Corrupt);
}
}
let record = decode_record(&payload)?;
frames.push(FrameIndex {
t: record.t,
offset: durable_len,
len: frame_len,
});
durable_len = durable_len
.checked_add(frame_len)
.ok_or(LogError::Corrupt)?;
}
Ok((frames, durable_len))
}
pub fn append_framed_record(out: &mut Vec<u8>, record: &TxRecord) -> Result<(), LogError> {
let payload = encode_record(record);
let header = frame_header(payload.len())?;
out.extend_from_slice(&header);
out.extend_from_slice(&payload);
out.extend_from_slice(&frame_checksum(header, &payload).to_be_bytes());
Ok(())
}
pub fn decode_framed_records(mut bytes: &[u8]) -> Result<Vec<TxRecord>, LogError> {
let mut records = Vec::new();
while !bytes.is_empty() {
if bytes.len() < 8 {
return Err(LogError::Corrupt);
}
let header: [u8; 8] = bytes[..8].try_into().map_err(|_| LogError::Corrupt)?;
let (payload_len, checksummed) = frame_payload_len(header)?;
bytes = &bytes[8..];
let payload = bytes.get(..payload_len).ok_or(LogError::Corrupt)?;
bytes = &bytes[payload_len..];
if checksummed {
let stored_checksum = u32::from_be_bytes(
bytes
.get(..FRAME_CHECKSUM_LEN)
.ok_or(LogError::Corrupt)?
.try_into()
.map_err(|_| LogError::Corrupt)?,
);
if stored_checksum != frame_checksum(header, payload) {
return Err(LogError::Corrupt);
}
bytes = &bytes[FRAME_CHECKSUM_LEN..];
}
records.push(decode_record(payload)?);
}
Ok(records)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn versioned_log_bounds_cached_read_descriptors() {
let dir = tempfile::tempdir().expect("tempdir");
let segment_count = MAX_CACHED_READ_VERSION_FILES + 5;
for version in 1..=u64::try_from(segment_count).expect("segment count fits u64") {
File::create(version_path(dir.path(), "db", version)).expect("create segment");
}
let log = VersionedLog::open_read_only(dir.path(), "db").expect("open log");
let state = log.state.read().expect("log lock");
assert_eq!(state.files.len(), segment_count);
assert!(
state
.files
.iter()
.filter(|file| file.file.file.is_some())
.count()
<= MAX_CACHED_READ_VERSION_FILES
);
}
}