use std::collections::HashMap;
use std::fs;
use std::io;
use gnitz_wire::PkBuf;
use gnitz_wire::MAX_PK_BYTES;
use gnitz_wire::{Reader, Writer};
use gnitz_zset::repr::StorageError;
const MAGIC: u64 = 0x4D414E49464E5447;
const VERSION: u64 = 17;
const MANIFEST_FILE: &str = "manifest.bin";
const STAGING_FILE: &str = "manifest.bin.tmp";
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(crate) struct ShardSet {
pub run_bytes: u64,
pub entries: Vec<ManifestEntry>,
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(crate) struct Manifest {
pub checkpoint_mark: u64,
pub caller_record: Vec<u8>,
pub shards: ShardSet,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct ManifestEntry {
pub seq: u64,
pub newest: u64,
pub level: u64,
pub guard_key: PkBuf,
}
pub(crate) fn encode(m: &Manifest) -> Vec<u8> {
let mut w = Writer::new();
w.u64(MAGIC)
.u64(VERSION)
.u64(m.checkpoint_mark)
.u64(m.shards.run_bytes)
.bytes32(&m.caller_record);
for e in &m.shards.entries {
w.u64(e.seq).u64(e.newest).u64(e.level).bytes32(e.guard_key.pk_bytes());
}
let mut buf = w.into_vec();
let digest = gnitz_wire::checksum(&buf);
buf.extend_from_slice(&digest.to_le_bytes());
buf
}
fn decode(buf: &[u8]) -> Result<Manifest, StorageError> {
const TRUNCATED: StorageError = StorageError::Corrupt("manifest truncated");
let truncated = |_| TRUNCATED;
let (covered, digest) = buf.split_last_chunk::<8>().ok_or(TRUNCATED)?;
let mut r = Reader::new(covered);
if r.u64().map_err(truncated)? != MAGIC {
return Err(StorageError::Corrupt("manifest magic"));
}
if r.u64().map_err(truncated)? != VERSION {
return Err(StorageError::Corrupt("manifest version"));
}
if gnitz_wire::checksum(covered) != u64::from_le_bytes(*digest) {
return Err(StorageError::Corrupt("manifest checksum"));
}
decode_body(&mut r).map_err(|_| StorageError::Corrupt("manifest body"))
}
fn decode_body(r: &mut Reader) -> Result<Manifest, String> {
let checkpoint_mark = r.u64()?;
let run_bytes = r.u64()?;
let caller_record = r.bytes32()?.to_vec();
let mut m = Manifest {
checkpoint_mark,
caller_record,
shards: ShardSet { run_bytes, entries: Vec::new() },
};
while r.remaining() > 0 {
let (seq, newest, level) = (r.u64()?, r.u64()?, r.u64()?);
let key = r.bytes32()?;
if key.len() > MAX_PK_BYTES {
return Err("guard key too wide".into());
}
m.shards.entries.push(ManifestEntry {
seq,
newest,
level,
guard_key: PkBuf::from_bytes(key),
});
}
Ok(m)
}
pub(crate) fn manifest_path(store_dir: &str) -> String {
format!("{store_dir}/{MANIFEST_FILE}")
}
pub(super) fn staging_path(dir: &str) -> String {
format!("{dir}/{STAGING_FILE}")
}
pub(super) fn shard_name(seq: u64) -> String {
format!("shard_{seq}.db")
}
fn shard_seq(name: &str) -> Option<u64> {
let seq = name.strip_prefix("shard_")?.strip_suffix(".db")?.parse().ok()?;
(shard_name(seq) == name).then_some(seq)
}
pub(super) fn shard_seqs(dir: &str) -> io::Result<Vec<u64>> {
let mut seqs = Vec::new();
for entry in fs::read_dir(dir)? {
seqs.extend(entry?.file_name().to_str().and_then(shard_seq));
}
Ok(seqs)
}
pub(crate) fn shard_path(dir: &str, seq: u64) -> String {
format!("{dir}/{}", shard_name(seq))
}
fn absent_ok(r: io::Result<()>) -> io::Result<()> {
match r {
Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(()),
r => r,
}
}
pub(super) fn remove_stale_files(dir: &str, live: impl Fn(u64) -> bool) -> io::Result<()> {
absent_ok(fs::remove_file(staging_path(dir)))?;
for seq in shard_seqs(dir)? {
if !live(seq) {
absent_ok(fs::remove_file(shard_path(dir, seq)))?;
}
}
Ok(())
}
pub(crate) fn shard_files(dir: &str) -> Result<Vec<(String, Option<u64>)>, StorageError> {
let entries = read_intact(dir)?.map_or(Vec::new(), |m| m.shards.entries);
let levels: HashMap<u64, u64> = entries.iter().map(|e| (e.seq, e.level)).collect();
let seqs = shard_seqs(dir)?.into_iter();
Ok(seqs
.map(|seq| (shard_path(dir, seq), levels.get(&seq).copied()))
.collect())
}
pub(crate) fn read(dir: &str) -> Result<Option<Manifest>, StorageError> {
match fs::read(manifest_path(dir)) {
Ok(buf) => decode(&buf).map(Some),
Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(None),
Err(e) => Err(e.into()),
}
}
pub(crate) fn read_intact(dir: &str) -> Result<Option<Manifest>, StorageError> {
match read(dir) {
Err(StorageError::Corrupt(_)) => Ok(None),
r => r,
}
}
pub(crate) fn read_at(dir: &str, generation: u64) -> Result<Option<Manifest>, StorageError> {
Ok(read_intact(dir)?.filter(|m| m.checkpoint_mark == generation))
}
pub(crate) fn fsync_dir(dir: &str) -> io::Result<()> {
fs::File::open(dir)?.sync_all()
}
pub(crate) fn prepare(dir: &str, bytes: &[u8]) -> Result<(), StorageError> {
fs::write(staging_path(dir), bytes)?;
Ok(())
}
pub(crate) fn commit(dir: &str) -> Result<(), StorageError> {
fs::rename(staging_path(dir), manifest_path(dir))?;
Ok(())
}
pub(crate) fn unlink(dir: &str) -> Result<(), StorageError> {
absent_ok(fs::remove_file(manifest_path(dir)))?;
Ok(absent_ok(fsync_dir(dir))?)
}
pub(crate) fn retire_store(store_dir: &str) -> Result<(), StorageError> {
unlink(store_dir)?;
if let Err(e) = absent_ok(fs::remove_dir_all(store_dir)) {
gnitz_warn!("storage: failed to remove retired store dir {}: {}", store_dir, e);
}
Ok(())
}
pub(crate) fn link_store(src_dir: &str, dst_dir: &str) -> Result<(), StorageError> {
let m = read(src_dir)?.ok_or(StorageError::Io(libc::ENOENT))?;
fs::create_dir_all(dst_dir)?;
for e in &m.shards.entries {
fs::hard_link(shard_path(src_dir, e.seq), shard_path(dst_dir, e.seq))?;
}
fsync_dir(dst_dir)?;
fs::hard_link(manifest_path(src_dir), manifest_path(dst_dir))?;
Ok(fsync_dir(dst_dir)?)
}
#[cfg(test)]
#[path = "tests/manifest.rs"]
mod tests;