use std::fs;
use std::io::Write;
use std::path::PathBuf;
use crate::error::Result;
use crate::segment::{filename, parse_filename, SegmentRange};
use crate::DurabilityPolicy;
const TMP_SUFFIX: &str = ".tmp";
pub trait SegmentStore: Send + Sync {
fn create_dir_all(&self) -> Result<()>;
fn scan(&self) -> Result<Vec<SegmentRange>>;
fn clean_tmp(&self) -> Result<usize>;
fn segment_size(&self, range: SegmentRange) -> u64;
fn remove_segment(&self, range: SegmentRange) -> Result<bool>;
fn write_atomic(
&self,
range: SegmentRange,
payload: &[u8],
policy: DurabilityPolicy,
) -> Result<u64>;
fn read_bytes(&self, range: SegmentRange) -> Result<Vec<u8>>;
}
#[derive(Debug)]
pub struct RealStore {
dir: PathBuf,
}
impl RealStore {
pub(crate) fn new(dir: PathBuf) -> Self {
Self { dir }
}
fn segment_path(&self, range: SegmentRange) -> PathBuf {
self.dir.join(filename(range.start, range.end))
}
}
impl SegmentStore for RealStore {
fn create_dir_all(&self) -> Result<()> {
Ok(fs::create_dir_all(&self.dir)?)
}
fn scan(&self) -> Result<Vec<SegmentRange>> {
let mut segments = Vec::new();
for entry in fs::read_dir(&self.dir)? {
let entry = entry?;
if let Some(range) = parse_filename(&entry.file_name().to_string_lossy()) {
segments.push(range);
}
}
segments.sort_by_key(|s| s.start);
Ok(segments)
}
fn clean_tmp(&self) -> Result<usize> {
let mut removed = 0usize;
for entry in fs::read_dir(&self.dir)? {
let entry = entry?;
let path = entry.path();
if path
.file_name()
.is_some_and(|n| n.to_string_lossy().ends_with(TMP_SUFFIX))
&& fs::remove_file(&path).is_ok()
{
removed += 1;
}
}
Ok(removed)
}
fn segment_size(&self, range: SegmentRange) -> u64 {
fs::metadata(self.segment_path(range))
.map(|m| m.len())
.unwrap_or(0)
}
fn remove_segment(&self, range: SegmentRange) -> Result<bool> {
match fs::remove_file(self.segment_path(range)) {
Ok(()) => Ok(true),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(false),
Err(e) => Err(e.into()),
}
}
fn write_atomic(
&self,
range: SegmentRange,
payload: &[u8],
policy: DurabilityPolicy,
) -> Result<u64> {
let seg_name = filename(range.start, range.end);
let seg_path = self.dir.join(&seg_name);
let tmp_path = self.dir.join(format!("{seg_name}{TMP_SUFFIX}"));
{
let mut file = fs::File::create(&tmp_path)?;
file.write_all(payload)?;
if !matches!(policy, DurabilityPolicy::Throughput) {
file.sync_all()?;
}
}
fs::rename(&tmp_path, &seg_path)?;
if matches!(policy, DurabilityPolicy::Maximal) {
let dir_file = fs::File::open(&self.dir)?;
dir_file.sync_all()?;
}
Ok(payload.len() as u64)
}
fn read_bytes(&self, range: SegmentRange) -> Result<Vec<u8>> {
let path = self.segment_path(range);
Ok(fs::read(&path)?)
}
}