use std::fs;
use std::io::Write;
use std::path::Path;
use serde::de::DeserializeOwned;
use serde::Serialize;
use crate::cipher::SegmentCipher;
use crate::error::{Result, SegmentError};
const SEGMENT_PREFIX: &str = "seg_";
const SEGMENT_SUFFIX: &str = ".zst";
const TMP_SUFFIX: &str = ".tmp";
const NONCE_LEN: usize = 12;
#[derive(Debug, Clone, Copy)]
pub(crate) struct SegmentRange {
pub(crate) start: u64,
pub(crate) end: u64,
}
pub(crate) fn filename(start: u64, end: u64) -> String {
format!("{SEGMENT_PREFIX}{start:012}_{end:012}{SEGMENT_SUFFIX}")
}
pub(crate) fn parse_filename(name: &str) -> Option<SegmentRange> {
let core = name
.strip_prefix(SEGMENT_PREFIX)?
.strip_suffix(SEGMENT_SUFFIX)?;
let (start_str, end_str) = core.split_once('_')?;
let start = start_str.parse().ok()?;
let end = end_str.parse().ok()?;
Some(SegmentRange { start, end })
}
pub(crate) fn scan(dir: &Path) -> Result<Vec<SegmentRange>> {
let mut segments = Vec::new();
for entry in fs::read_dir(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)
}
pub(crate) fn clean_tmp(dir: &Path) -> Result<()> {
for entry in fs::read_dir(dir)? {
let entry = entry?;
let path = entry.path();
if path
.file_name()
.is_some_and(|n| n.to_string_lossy().ends_with(TMP_SUFFIX))
{
let _ = fs::remove_file(&path);
}
}
Ok(())
}
fn encode<T: Serialize>(
cipher: Option<&dyn SegmentCipher>,
level: i32,
events: &[T],
) -> Result<Vec<u8>> {
let mut cbor_buf = Vec::new();
ciborium::into_writer(events, &mut cbor_buf)
.map_err(|e| SegmentError::Cbor(format!("serialization: {e}")))?;
let compressed = zstd::encode_all(cbor_buf.as_slice(), level)?;
match cipher {
Some(cipher) => cipher.encrypt(&compressed),
None => Ok(compressed),
}
}
fn decode<T: DeserializeOwned>(cipher: Option<&dyn SegmentCipher>, raw: Vec<u8>) -> Result<Vec<T>> {
let compressed = match cipher {
Some(cipher) => cipher.decrypt(&raw)?,
None => raw,
};
let cbor_buf = zstd::decode_all(compressed.as_slice())?;
ciborium::from_reader(cbor_buf.as_slice())
.map_err(|e| SegmentError::Cbor(format!("deserialization: {e}")))
}
pub(crate) fn write<T: Serialize>(
dir: &Path,
cipher: Option<&dyn SegmentCipher>,
level: i32,
range: SegmentRange,
events: &[T],
) -> Result<u64> {
let final_bytes = encode(cipher, level, events)?;
let seg_name = filename(range.start, range.end);
let seg_path = dir.join(&seg_name);
let tmp_path = dir.join(format!("{seg_name}{TMP_SUFFIX}"));
{
let mut file = fs::File::create(&tmp_path)?;
file.write_all(&final_bytes)?;
file.sync_all()?;
}
fs::rename(&tmp_path, &seg_path)?;
Ok(final_bytes.len() as u64)
}
pub(crate) fn read<T: DeserializeOwned>(
dir: &Path,
cipher: Option<&dyn SegmentCipher>,
range: SegmentRange,
) -> Result<Vec<T>> {
let path = dir.join(filename(range.start, range.end));
let raw = fs::read(&path)?;
if cipher.is_some() && raw.len() < NONCE_LEN {
return Err(SegmentError::Integrity(format!(
"segment {} too small for nonce ({} bytes)",
path.display(),
raw.len()
)));
}
decode(cipher, raw)
}