use super::error::{Error, Result};
use serde::{Deserialize, Serialize};
use std::fs;
use std::path::{Path, PathBuf};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ReplicationState {
pub device: u64,
pub inode: u64,
pub db_path: PathBuf,
pub base_id: Option<String>,
pub last_committed_offset: u64,
pub last_uploaded_etag: Option<String>,
pub last_seen_size: u64,
}
impl ReplicationState {
pub fn load_or_create(state_path: &Path, db_path: &Path) -> Result<Self> {
if state_path.exists() {
let content = fs::read_to_string(state_path)?;
let state: ReplicationState = toml::from_str(&content)
.map_err(|e| Error::InvalidState(format!("Failed to parse state file: {}", e)))?;
let current_meta = Self::get_file_metadata(db_path)?;
if state.device != current_meta.0 || state.inode != current_meta.1 {
return Ok(Self::new(db_path));
}
Ok(state)
} else {
Ok(Self::new(db_path))
}
}
pub fn new(db_path: &Path) -> Self {
let (device, inode) = Self::get_file_metadata(db_path).unwrap_or((0, 0));
let last_seen_size = fs::metadata(db_path).map(|m| m.len()).unwrap_or(0);
Self {
device,
inode,
db_path: db_path.to_path_buf(),
base_id: None,
last_committed_offset: 64, last_uploaded_etag: None,
last_seen_size,
}
}
#[cfg(unix)]
fn get_file_metadata(path: &Path) -> Result<(u64, u64)> {
use std::os::unix::fs::MetadataExt;
let meta = fs::metadata(path)?;
Ok((meta.dev(), meta.ino()))
}
#[cfg(not(unix))]
fn get_file_metadata(path: &Path) -> Result<(u64, u64)> {
let meta = fs::metadata(path)?;
let modified = meta.modified()?;
let mtime = modified
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
path.hash(&mut hasher);
mtime.hash(&mut hasher);
Ok((hasher.finish(), mtime))
}
pub fn save(&self, state_path: &Path) -> Result<()> {
let content = toml::to_string_pretty(self)
.map_err(|e| Error::InvalidState(format!("Failed to serialize state: {}", e)))?;
let temp_path = state_path.with_extension("tmp");
fs::write(&temp_path, content)?;
fs::rename(&temp_path, state_path)?;
Ok(())
}
pub fn check_rotation(&mut self) -> Result<bool> {
let (device, inode) = Self::get_file_metadata(&self.db_path)?;
let current_size = fs::metadata(&self.db_path)?.len();
if device != self.device || inode != self.inode || current_size < self.last_seen_size {
self.device = device;
self.inode = inode;
self.base_id = None;
self.last_committed_offset = 64; self.last_uploaded_etag = None;
self.last_seen_size = current_size;
return Ok(true);
}
self.last_seen_size = current_size;
Ok(false)
}
pub fn update_after_upload(&mut self, offset: u64, etag: Option<String>) {
self.last_committed_offset = offset;
self.last_uploaded_etag = etag;
}
pub fn update_after_base(&mut self, base_id: String, offset: u64) {
self.base_id = Some(base_id);
self.last_committed_offset = offset;
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_state_serialization() {
let db_path = PathBuf::from("/tmp/test.teg");
let state = ReplicationState::new(&db_path);
let state_str = toml::to_string(&state).unwrap();
let _parsed: ReplicationState = toml::from_str(&state_str).unwrap();
}
}