use std::fs::{File, OpenOptions};
use std::io::{self, BufWriter, Read, Write};
use std::path::{Path, PathBuf};
const OP_PUT: u8 = 0x00;
const OP_DELETE: u8 = 0x01;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct WalEntry {
pub key: String,
pub value: Option<Vec<u8>>,
}
pub struct WriteAheadLog {
path: PathBuf,
writer: Option<BufWriter<File>>,
}
impl WriteAheadLog {
pub fn open(path: impl AsRef<Path>) -> io::Result<Self> {
let path = path.as_ref().to_path_buf();
let file = OpenOptions::new()
.read(true)
.create(true)
.append(true)
.open(&path)?;
Ok(Self {
path,
writer: Some(BufWriter::new(file)),
})
}
fn writer(&mut self) -> &mut BufWriter<File> {
self.writer
.as_mut()
.expect("wal writer dropped without reopen")
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn log_put(&mut self, key: &str, value: &[u8]) -> io::Result<()> {
self.append(OP_PUT, key.as_bytes(), value)
}
pub fn log_delete(&mut self, key: &str) -> io::Result<()> {
self.append(OP_DELETE, key.as_bytes(), &[])
}
fn append(&mut self, op: u8, key: &[u8], value: &[u8]) -> io::Result<()> {
let mut crc = Crc32::new();
crc.update(&[op]);
let kl = key.len() as u32;
let vl = value.len() as u32;
crc.update(key);
crc.update(value);
let checksum = crc.finalize();
let w = self.writer();
w.write_all(&[op])?;
w.write_all(&kl.to_be_bytes())?;
w.write_all(key)?;
w.write_all(&vl.to_be_bytes())?;
w.write_all(value)?;
w.write_all(&checksum.to_be_bytes())?;
w.flush()
}
pub fn sync(&mut self) -> io::Result<()> {
let w = self.writer();
w.flush()?;
w.get_ref().sync_all()
}
pub fn truncate(&mut self) -> io::Result<()> {
if let Some(mut w) = self.writer.take() {
w.flush()?;
drop(w);
}
let _ = OpenOptions::new()
.write(true)
.truncate(true)
.create(true)
.open(&self.path)?;
let file = OpenOptions::new()
.read(true)
.create(true)
.append(true)
.open(&self.path)?;
self.writer = Some(BufWriter::new(file));
Ok(())
}
pub fn replay(path: impl AsRef<Path>) -> io::Result<Vec<WalEntry>> {
let path = path.as_ref();
if !path.exists() {
return Ok(Vec::new());
}
let mut buf = Vec::new();
File::open(path)?.read_to_end(&mut buf)?;
let mut entries = Vec::new();
let mut p = 0usize;
while p < buf.len() {
if p + 5 > buf.len() {
break;
}
let op = buf[p];
let key_len = u32::from_be_bytes(buf[p + 1..p + 5].try_into().unwrap()) as usize;
let after_key = p + 5 + key_len;
if after_key + 4 > buf.len() {
break;
}
let value_len =
u32::from_be_bytes(buf[after_key..after_key + 4].try_into().unwrap()) as usize;
let after_value = after_key + 4 + value_len;
if after_value + 4 > buf.len() {
break;
}
let key = &buf[p + 5..p + 5 + key_len];
let value = &buf[after_key + 4..after_key + 4 + value_len];
let stored_crc =
u32::from_be_bytes(buf[after_value..after_value + 4].try_into().unwrap());
let mut crc = Crc32::new();
crc.update(&[op]);
crc.update(key);
crc.update(value);
if crc.finalize() != stored_crc {
break;
}
let key_str = match std::str::from_utf8(key) {
Ok(s) => s.to_string(),
Err(_) => break,
};
let value_opt = match op {
OP_PUT => Some(value.to_vec()),
OP_DELETE => None,
_ => break,
};
entries.push(WalEntry {
key: key_str,
value: value_opt,
});
p = after_value + 4;
}
Ok(entries)
}
}
struct Crc32 {
state: u32,
}
impl Crc32 {
fn new() -> Self {
Self { state: 0xffff_ffff }
}
fn update(&mut self, bytes: &[u8]) {
let mut s = self.state;
for &b in bytes {
let i = ((s ^ b as u32) & 0xff) as usize;
s = (s >> 8) ^ CRC32_TABLE[i];
}
self.state = s;
}
fn finalize(self) -> u32 {
self.state ^ 0xffff_ffff
}
}
const CRC32_POLY: u32 = 0xedb8_8320;
static CRC32_TABLE: [u32; 256] = {
let mut table = [0u32; 256];
let mut i = 0;
while i < 256 {
let mut c = i as u32;
let mut j = 0;
while j < 8 {
c = if c & 1 != 0 {
(c >> 1) ^ CRC32_POLY
} else {
c >> 1
};
j += 1;
}
table[i] = c;
i += 1;
}
table
};
#[cfg(test)]
#[path = "wal_tests.rs"]
mod tests;