use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, HashMap};
use std::fmt;
use std::path::{Path, PathBuf};
use std::time::{Duration, SystemTime};
use crate::error::{Result, TdbError};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum WalRecordType {
Insert,
Delete,
TxnBegin,
TxnCommit,
TxnAbort,
Checkpoint,
SchemaChange,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WalRecord {
pub lsn: u64,
pub record_type: WalRecordType,
pub txn_id: u64,
pub timestamp: u64,
pub key: Vec<u8>,
pub value: Vec<u8>,
pub checksum: u32,
}
impl WalRecord {
pub fn new(
lsn: u64,
record_type: WalRecordType,
txn_id: u64,
key: Vec<u8>,
value: Vec<u8>,
) -> Self {
let timestamp = SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
let checksum = Self::compute_checksum(&key, &value);
Self {
lsn,
record_type,
txn_id,
timestamp,
key,
value,
checksum,
}
}
fn compute_checksum(key: &[u8], value: &[u8]) -> u32 {
let mut hasher = crc32fast::Hasher::new();
hasher.update(key);
hasher.update(value);
hasher.finalize()
}
pub fn verify_checksum(&self) -> bool {
let expected = Self::compute_checksum(&self.key, &self.value);
self.checksum == expected
}
pub fn estimated_size(&self) -> usize {
8 + 1 + 8 + 8 + self.key.len() + self.value.len() + 4
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WalSegment {
pub segment_id: u64,
pub records: Vec<WalRecord>,
pub first_lsn: u64,
pub last_lsn: u64,
pub compacted: bool,
pub created_at: u64,
}
impl WalSegment {
pub fn new(segment_id: u64) -> Self {
let now = SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
Self {
segment_id,
records: Vec::new(),
first_lsn: 0,
last_lsn: 0,
compacted: false,
created_at: now,
}
}
pub fn add_record(&mut self, record: WalRecord) {
if self.records.is_empty() {
self.first_lsn = record.lsn;
}
self.last_lsn = record.lsn;
self.records.push(record);
}
pub fn len(&self) -> usize {
self.records.len()
}
pub fn is_empty(&self) -> bool {
self.records.is_empty()
}
pub fn estimated_size(&self) -> usize {
self.records.iter().map(|r| r.estimated_size()).sum()
}
pub fn timestamp_range(&self) -> Option<(u64, u64)> {
if self.records.is_empty() {
return None;
}
let first_ts = self.records.first().map(|r| r.timestamp).unwrap_or(0);
let last_ts = self.records.last().map(|r| r.timestamp).unwrap_or(0);
Some((first_ts, last_ts))
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CompactionConfig {
pub min_segments_to_compact: usize,
pub max_segments_per_pass: usize,
pub remove_dead_entries: bool,
pub preserve_transactions: bool,
pub verify_checksums: bool,
}
impl Default for CompactionConfig {
fn default() -> Self {
Self {
min_segments_to_compact: 3,
max_segments_per_pass: 10,
remove_dead_entries: true,
preserve_transactions: true,
verify_checksums: true,
}
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct CompactionStats {
pub segments_compacted: usize,
pub records_before: usize,
pub records_after: usize,
pub dead_entries_removed: usize,
pub bytes_saved: usize,
pub duration_ms: u64,
pub checksum_failures: usize,
}
pub struct WalCompactor {
config: CompactionConfig,
}
impl WalCompactor {
pub fn new(config: CompactionConfig) -> Self {
Self { config }
}
pub fn compact(&self, segments: &[WalSegment]) -> Result<(WalSegment, CompactionStats)> {
let start = std::time::Instant::now();
let mut stats = CompactionStats::default();
if segments.is_empty() {
return Ok((WalSegment::new(0), stats));
}
stats.segments_compacted = segments.len();
let mut all_records: Vec<&WalRecord> =
segments.iter().flat_map(|s| s.records.iter()).collect();
all_records.sort_by_key(|r| r.lsn);
stats.records_before = all_records.len();
if self.config.verify_checksums {
for record in &all_records {
if !record.verify_checksum() {
stats.checksum_failures += 1;
}
}
}
let mut key_state: HashMap<Vec<u8>, WalRecordType> = HashMap::new();
let mut live_lsns: std::collections::HashSet<u64> = std::collections::HashSet::new();
if self.config.remove_dead_entries {
for record in &all_records {
match record.record_type {
WalRecordType::Insert => {
key_state.insert(record.key.clone(), WalRecordType::Insert);
}
WalRecordType::Delete => {
key_state.insert(record.key.clone(), WalRecordType::Delete);
}
_ => {}
}
}
for record in &all_records {
let keep = match record.record_type {
WalRecordType::Insert => {
matches!(key_state.get(&record.key), Some(WalRecordType::Insert))
}
WalRecordType::Delete => {
false
}
WalRecordType::TxnBegin
| WalRecordType::TxnCommit
| WalRecordType::TxnAbort => self.config.preserve_transactions,
WalRecordType::Checkpoint | WalRecordType::SchemaChange => true,
};
if keep {
live_lsns.insert(record.lsn);
}
}
} else {
for record in &all_records {
live_lsns.insert(record.lsn);
}
}
let first_segment_id = segments.first().map(|s| s.segment_id).unwrap_or(0);
let mut compacted = WalSegment::new(first_segment_id);
compacted.compacted = true;
let size_before: usize = all_records.iter().map(|r| r.estimated_size()).sum();
for record in &all_records {
if live_lsns.contains(&record.lsn) {
compacted.add_record((*record).clone());
}
}
stats.records_after = compacted.len();
stats.dead_entries_removed = stats.records_before - stats.records_after;
let size_after: usize = compacted.records.iter().map(|r| r.estimated_size()).sum();
stats.bytes_saved = size_before.saturating_sub(size_after);
stats.duration_ms = start.elapsed().as_millis() as u64;
Ok((compacted, stats))
}
pub fn should_compact(&self, segment_count: usize) -> bool {
segment_count >= self.config.min_segments_to_compact
}
pub fn config(&self) -> &CompactionConfig {
&self.config
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BaseSnapshot {
pub snapshot_id: u64,
pub lsn: u64,
pub timestamp: u64,
pub path: PathBuf,
pub size_bytes: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PitrConfig {
pub enabled: bool,
pub snapshot_interval_records: u64,
pub max_snapshots: usize,
pub verify_checksums: bool,
}
impl Default for PitrConfig {
fn default() -> Self {
Self {
enabled: true,
snapshot_interval_records: 100_000,
max_snapshots: 5,
verify_checksums: true,
}
}
}
#[derive(Debug, Clone)]
pub struct RecoveryPlan {
pub base_snapshot: Option<BaseSnapshot>,
pub segments_to_replay: Vec<u64>,
pub target_lsn: Option<u64>,
pub target_timestamp: Option<u64>,
pub estimated_records: u64,
}
impl RecoveryPlan {
pub fn needs_snapshot(&self) -> bool {
self.base_snapshot.is_some()
}
pub fn needs_wal_replay(&self) -> bool {
!self.segments_to_replay.is_empty()
}
}
pub struct PitrEngine {
config: PitrConfig,
snapshots: BTreeMap<u64, BaseSnapshot>,
segments: BTreeMap<u64, WalSegment>,
}
impl PitrEngine {
pub fn new(config: PitrConfig) -> Self {
Self {
config,
snapshots: BTreeMap::new(),
segments: BTreeMap::new(),
}
}
pub fn register_snapshot(&mut self, snapshot: BaseSnapshot) {
self.snapshots.insert(snapshot.lsn, snapshot);
while self.snapshots.len() > self.config.max_snapshots {
if let Some((&oldest_lsn, _)) = self.snapshots.iter().next() {
self.snapshots.remove(&oldest_lsn);
}
}
}
pub fn register_segment(&mut self, segment: WalSegment) {
self.segments.insert(segment.segment_id, segment);
}
pub fn plan_recovery_to_lsn(&self, target_lsn: u64) -> Result<RecoveryPlan> {
let base_snapshot = self
.snapshots
.range(..=target_lsn)
.next_back()
.map(|(_, s)| s.clone());
let start_lsn = base_snapshot.as_ref().map(|s| s.lsn).unwrap_or(0);
let segments_to_replay: Vec<u64> = self
.segments
.values()
.filter(|s| s.last_lsn > start_lsn && s.first_lsn <= target_lsn)
.map(|s| s.segment_id)
.collect();
let estimated_records: u64 = self
.segments
.values()
.filter(|s| segments_to_replay.contains(&s.segment_id))
.map(|s| s.records.len() as u64)
.sum();
Ok(RecoveryPlan {
base_snapshot,
segments_to_replay,
target_lsn: Some(target_lsn),
target_timestamp: None,
estimated_records,
})
}
pub fn plan_recovery_to_timestamp(&self, target_timestamp: u64) -> Result<RecoveryPlan> {
let base_snapshot = self
.snapshots
.values()
.rfind(|s| s.timestamp <= target_timestamp)
.cloned();
let start_lsn = base_snapshot.as_ref().map(|s| s.lsn).unwrap_or(0);
let segments_to_replay: Vec<u64> = self
.segments
.values()
.filter(|s| {
if let Some((_, last_ts)) = s.timestamp_range() {
last_ts
> base_snapshot
.as_ref()
.map(|snap| snap.timestamp)
.unwrap_or(0)
&& s.first_lsn > start_lsn
} else {
false
}
})
.map(|s| s.segment_id)
.collect();
let estimated_records: u64 = self
.segments
.values()
.filter(|s| segments_to_replay.contains(&s.segment_id))
.map(|s| s.records.len() as u64)
.sum();
Ok(RecoveryPlan {
base_snapshot,
segments_to_replay,
target_lsn: None,
target_timestamp: Some(target_timestamp),
estimated_records,
})
}
pub fn execute_plan(&self, plan: &RecoveryPlan) -> Result<Vec<WalRecord>> {
let mut recovered_records = Vec::new();
for seg_id in &plan.segments_to_replay {
if let Some(segment) = self.segments.get(seg_id) {
for record in &segment.records {
if let Some(target_lsn) = plan.target_lsn {
if record.lsn > target_lsn {
break;
}
}
if let Some(target_ts) = plan.target_timestamp {
if record.timestamp > target_ts {
break;
}
}
if self.config.verify_checksums && !record.verify_checksum() {
return Err(TdbError::Wal(format!(
"Checksum verification failed for record LSN {}",
record.lsn
)));
}
recovered_records.push(record.clone());
}
}
}
Ok(recovered_records)
}
pub fn snapshot_count(&self) -> usize {
self.snapshots.len()
}
pub fn segment_count(&self) -> usize {
self.segments.len()
}
pub fn latest_snapshot_lsn(&self) -> Option<u64> {
self.snapshots.keys().next_back().copied()
}
pub fn config(&self) -> &PitrConfig {
&self.config
}
pub fn should_take_snapshot(&self, current_lsn: u64) -> bool {
if !self.config.enabled {
return false;
}
let latest_lsn = self.latest_snapshot_lsn().unwrap_or(0);
current_lsn.saturating_sub(latest_lsn) >= self.config.snapshot_interval_records
}
}
#[cfg(test)]
mod tests {
use super::*;
fn make_record(lsn: u64, rt: WalRecordType, key: &[u8], value: &[u8]) -> WalRecord {
WalRecord::new(lsn, rt, 1, key.to_vec(), value.to_vec())
}
fn make_segment(id: u64, records: Vec<WalRecord>) -> WalSegment {
let mut seg = WalSegment::new(id);
for r in records {
seg.add_record(r);
}
seg
}
#[test]
fn test_wal_record_creation() {
let rec = make_record(1, WalRecordType::Insert, b"key1", b"val1");
assert_eq!(rec.lsn, 1);
assert_eq!(rec.record_type, WalRecordType::Insert);
assert_eq!(rec.key, b"key1");
assert_eq!(rec.value, b"val1");
}
#[test]
fn test_wal_record_checksum_valid() {
let rec = make_record(1, WalRecordType::Insert, b"key", b"value");
assert!(rec.verify_checksum());
}
#[test]
fn test_wal_record_checksum_invalid() {
let mut rec = make_record(1, WalRecordType::Insert, b"key", b"value");
rec.key = b"modified".to_vec(); assert!(!rec.verify_checksum());
}
#[test]
fn test_wal_record_estimated_size() {
let rec = make_record(1, WalRecordType::Insert, b"key", b"value");
assert!(rec.estimated_size() > 0);
}
#[test]
fn test_segment_new() {
let seg = WalSegment::new(1);
assert_eq!(seg.segment_id, 1);
assert!(seg.is_empty());
assert_eq!(seg.len(), 0);
}
#[test]
fn test_segment_add_records() {
let mut seg = WalSegment::new(1);
seg.add_record(make_record(1, WalRecordType::Insert, b"k1", b"v1"));
seg.add_record(make_record(2, WalRecordType::Insert, b"k2", b"v2"));
assert_eq!(seg.len(), 2);
assert_eq!(seg.first_lsn, 1);
assert_eq!(seg.last_lsn, 2);
}
#[test]
fn test_segment_estimated_size() {
let seg = make_segment(
1,
vec![
make_record(1, WalRecordType::Insert, b"k1", b"v1"),
make_record(2, WalRecordType::Insert, b"k2", b"v2"),
],
);
assert!(seg.estimated_size() > 0);
}
#[test]
fn test_segment_timestamp_range() {
let seg = make_segment(
1,
vec![
make_record(1, WalRecordType::Insert, b"k1", b"v1"),
make_record(2, WalRecordType::Insert, b"k2", b"v2"),
],
);
let range = seg.timestamp_range();
assert!(range.is_some());
let (first, last) = range.expect("should have range");
assert!(last >= first);
}
#[test]
fn test_segment_timestamp_range_empty() {
let seg = WalSegment::new(1);
assert!(seg.timestamp_range().is_none());
}
#[test]
fn test_compaction_config_default() {
let config = CompactionConfig::default();
assert_eq!(config.min_segments_to_compact, 3);
assert_eq!(config.max_segments_per_pass, 10);
assert!(config.remove_dead_entries);
assert!(config.preserve_transactions);
assert!(config.verify_checksums);
}
#[test]
fn test_compactor_should_compact() {
let compactor = WalCompactor::new(CompactionConfig::default());
assert!(!compactor.should_compact(1));
assert!(!compactor.should_compact(2));
assert!(compactor.should_compact(3));
assert!(compactor.should_compact(10));
}
#[test]
fn test_compact_empty_segments() {
let compactor = WalCompactor::new(CompactionConfig::default());
let (result, stats) = compactor.compact(&[]).expect("should succeed");
assert!(result.is_empty());
assert_eq!(stats.segments_compacted, 0);
}
#[test]
fn test_compact_single_segment() {
let compactor = WalCompactor::new(CompactionConfig::default());
let seg = make_segment(
1,
vec![
make_record(1, WalRecordType::Insert, b"k1", b"v1"),
make_record(2, WalRecordType::Insert, b"k2", b"v2"),
],
);
let (result, stats) = compactor.compact(&[seg]).expect("should succeed");
assert_eq!(stats.segments_compacted, 1);
assert_eq!(stats.records_before, 2);
assert_eq!(stats.records_after, 2);
}
#[test]
fn test_compact_removes_dead_entries() {
let compactor = WalCompactor::new(CompactionConfig::default());
let seg1 = make_segment(1, vec![make_record(1, WalRecordType::Insert, b"k1", b"v1")]);
let seg2 = make_segment(2, vec![make_record(2, WalRecordType::Delete, b"k1", b"")]);
let (result, stats) = compactor.compact(&[seg1, seg2]).expect("should succeed");
assert_eq!(stats.segments_compacted, 2);
assert_eq!(stats.records_before, 2);
assert_eq!(stats.dead_entries_removed, 2);
assert_eq!(stats.records_after, 0);
}
#[test]
fn test_compact_preserves_live_entries() {
let compactor = WalCompactor::new(CompactionConfig::default());
let seg1 = make_segment(
1,
vec![
make_record(1, WalRecordType::Insert, b"k1", b"v1"),
make_record(2, WalRecordType::Insert, b"k2", b"v2"),
],
);
let seg2 = make_segment(2, vec![make_record(3, WalRecordType::Delete, b"k1", b"")]);
let (result, stats) = compactor.compact(&[seg1, seg2]).expect("should succeed");
assert_eq!(stats.records_after, 1);
}
#[test]
fn test_compact_preserves_txn_records() {
let config = CompactionConfig {
preserve_transactions: true,
..Default::default()
};
let compactor = WalCompactor::new(config);
let seg = make_segment(
1,
vec![
make_record(1, WalRecordType::TxnBegin, b"", b""),
make_record(2, WalRecordType::Insert, b"k1", b"v1"),
make_record(3, WalRecordType::TxnCommit, b"", b""),
],
);
let (result, stats) = compactor.compact(&[seg]).expect("should succeed");
assert_eq!(stats.records_after, 3);
}
#[test]
fn test_compact_without_dead_entry_removal() {
let config = CompactionConfig {
remove_dead_entries: false,
..Default::default()
};
let compactor = WalCompactor::new(config);
let seg = make_segment(
1,
vec![
make_record(1, WalRecordType::Insert, b"k1", b"v1"),
make_record(2, WalRecordType::Delete, b"k1", b""),
],
);
let (_, stats) = compactor.compact(&[seg]).expect("should succeed");
assert_eq!(stats.records_after, 2); assert_eq!(stats.dead_entries_removed, 0);
}
#[test]
fn test_compact_checksum_failures() {
let compactor = WalCompactor::new(CompactionConfig::default());
let mut rec = make_record(1, WalRecordType::Insert, b"k1", b"v1");
rec.checksum = 0xDEADBEEF;
let seg = make_segment(1, vec![rec]);
let (_, stats) = compactor.compact(&[seg]).expect("should succeed");
assert_eq!(stats.checksum_failures, 1);
}
#[test]
fn test_compact_preserves_checkpoint() {
let compactor = WalCompactor::new(CompactionConfig::default());
let seg = make_segment(
1,
vec![make_record(1, WalRecordType::Checkpoint, b"snap", b"data")],
);
let (result, stats) = compactor.compact(&[seg]).expect("should succeed");
assert_eq!(stats.records_after, 1);
}
#[test]
fn test_compact_preserves_schema_change() {
let compactor = WalCompactor::new(CompactionConfig::default());
let seg = make_segment(
1,
vec![make_record(
1,
WalRecordType::SchemaChange,
b"schema",
b"data",
)],
);
let (_, stats) = compactor.compact(&[seg]).expect("should succeed");
assert_eq!(stats.records_after, 1);
}
#[test]
fn test_compact_stats_bytes_saved() {
let compactor = WalCompactor::new(CompactionConfig::default());
let seg1 = make_segment(
1,
vec![make_record(1, WalRecordType::Insert, b"key", b"val")],
);
let seg2 = make_segment(2, vec![make_record(2, WalRecordType::Delete, b"key", b"")]);
let (_, stats) = compactor.compact(&[seg1, seg2]).expect("should succeed");
assert!(stats.bytes_saved > 0);
}
#[test]
fn test_pitr_config_default() {
let config = PitrConfig::default();
assert!(config.enabled);
assert_eq!(config.snapshot_interval_records, 100_000);
assert_eq!(config.max_snapshots, 5);
}
#[test]
fn test_pitr_engine_creation() {
let engine = PitrEngine::new(PitrConfig::default());
assert_eq!(engine.snapshot_count(), 0);
assert_eq!(engine.segment_count(), 0);
}
#[test]
fn test_pitr_register_snapshot() {
let mut engine = PitrEngine::new(PitrConfig::default());
engine.register_snapshot(BaseSnapshot {
snapshot_id: 1,
lsn: 100,
timestamp: 1000,
path: std::env::temp_dir().join("oxirs_tdb_snap-1"),
size_bytes: 1024,
});
assert_eq!(engine.snapshot_count(), 1);
assert_eq!(engine.latest_snapshot_lsn(), Some(100));
}
#[test]
fn test_pitr_register_segment() {
let mut engine = PitrEngine::new(PitrConfig::default());
let seg = make_segment(1, vec![make_record(1, WalRecordType::Insert, b"k", b"v")]);
engine.register_segment(seg);
assert_eq!(engine.segment_count(), 1);
}
#[test]
fn test_pitr_max_snapshots_enforced() {
let config = PitrConfig {
max_snapshots: 2,
..Default::default()
};
let mut engine = PitrEngine::new(config);
for i in 0..5u64 {
engine.register_snapshot(BaseSnapshot {
snapshot_id: i,
lsn: i * 100,
timestamp: i * 1000,
path: std::env::temp_dir().join(format!("oxirs_tdb_snap-{i}")),
size_bytes: 1024,
});
}
assert_eq!(engine.snapshot_count(), 2);
}
#[test]
fn test_pitr_plan_recovery_to_lsn() {
let mut engine = PitrEngine::new(PitrConfig::default());
engine.register_snapshot(BaseSnapshot {
snapshot_id: 1,
lsn: 100,
timestamp: 1000,
path: std::env::temp_dir().join("oxirs_tdb_snap-1"),
size_bytes: 1024,
});
let seg = make_segment(
1,
vec![
make_record(101, WalRecordType::Insert, b"k1", b"v1"),
make_record(150, WalRecordType::Insert, b"k2", b"v2"),
make_record(200, WalRecordType::Insert, b"k3", b"v3"),
],
);
engine.register_segment(seg);
let plan = engine.plan_recovery_to_lsn(150).expect("should succeed");
assert!(plan.needs_snapshot());
assert!(plan.needs_wal_replay());
assert_eq!(plan.target_lsn, Some(150));
}
#[test]
fn test_pitr_plan_recovery_no_snapshot() {
let mut engine = PitrEngine::new(PitrConfig::default());
let seg = make_segment(1, vec![make_record(1, WalRecordType::Insert, b"k1", b"v1")]);
engine.register_segment(seg);
let plan = engine.plan_recovery_to_lsn(1).expect("should succeed");
assert!(!plan.needs_snapshot());
assert!(plan.needs_wal_replay());
}
#[test]
fn test_pitr_execute_plan_lsn_bounded() {
let mut engine = PitrEngine::new(PitrConfig::default());
let seg = make_segment(
1,
vec![
make_record(1, WalRecordType::Insert, b"k1", b"v1"),
make_record(2, WalRecordType::Insert, b"k2", b"v2"),
make_record(3, WalRecordType::Insert, b"k3", b"v3"),
],
);
engine.register_segment(seg);
let plan = engine.plan_recovery_to_lsn(2).expect("should succeed");
let records = engine.execute_plan(&plan).expect("should succeed");
assert_eq!(records.len(), 2);
assert_eq!(records[0].lsn, 1);
assert_eq!(records[1].lsn, 2);
}
#[test]
fn test_pitr_execute_plan_checksum_failure() {
let config = PitrConfig {
verify_checksums: true,
..Default::default()
};
let mut engine = PitrEngine::new(config);
let mut rec = make_record(1, WalRecordType::Insert, b"k", b"v");
rec.checksum = 0xBAD;
let seg = make_segment(1, vec![rec]);
engine.register_segment(seg);
let plan = engine.plan_recovery_to_lsn(1).expect("should succeed");
let result = engine.execute_plan(&plan);
assert!(result.is_err());
}
#[test]
fn test_pitr_should_take_snapshot() {
let config = PitrConfig {
snapshot_interval_records: 100,
..Default::default()
};
let mut engine = PitrEngine::new(config);
assert!(engine.should_take_snapshot(100));
assert!(!engine.should_take_snapshot(50));
engine.register_snapshot(BaseSnapshot {
snapshot_id: 1,
lsn: 50,
timestamp: 1000,
path: std::env::temp_dir().join("oxirs_tdb_snap"),
size_bytes: 0,
});
assert!(!engine.should_take_snapshot(100)); assert!(engine.should_take_snapshot(150)); }
#[test]
fn test_pitr_disabled() {
let config = PitrConfig {
enabled: false,
..Default::default()
};
let engine = PitrEngine::new(config);
assert!(!engine.should_take_snapshot(999999));
}
#[test]
fn test_recovery_plan_needs() {
let plan = RecoveryPlan {
base_snapshot: None,
segments_to_replay: vec![1],
target_lsn: Some(100),
target_timestamp: None,
estimated_records: 50,
};
assert!(!plan.needs_snapshot());
assert!(plan.needs_wal_replay());
let plan_with_snap = RecoveryPlan {
base_snapshot: Some(BaseSnapshot {
snapshot_id: 1,
lsn: 50,
timestamp: 1000,
path: std::env::temp_dir().join("oxirs_tdb_snap"),
size_bytes: 0,
}),
segments_to_replay: vec![],
target_lsn: Some(50),
target_timestamp: None,
estimated_records: 0,
};
assert!(plan_with_snap.needs_snapshot());
assert!(!plan_with_snap.needs_wal_replay());
}
#[test]
fn test_compactor_config_access() {
let config = CompactionConfig {
min_segments_to_compact: 5,
..Default::default()
};
let compactor = WalCompactor::new(config);
assert_eq!(compactor.config().min_segments_to_compact, 5);
}
#[test]
fn test_pitr_engine_config_access() {
let config = PitrConfig {
max_snapshots: 10,
..Default::default()
};
let engine = PitrEngine::new(config);
assert_eq!(engine.config().max_snapshots, 10);
}
}