1pub mod buffer_pool;
9pub mod checkpoint;
10pub mod free_space;
11pub mod memtable;
12pub mod metrics;
13pub mod sstable;
14pub mod wal;
15
16use std::collections::VecDeque;
17use std::collections::{BTreeMap, HashSet};
18use std::fs;
19use std::path::{Path, PathBuf};
20use std::sync::atomic::{AtomicU64, Ordering};
21use std::sync::{Arc, Mutex, RwLock};
22use std::time::{Instant, SystemTime, UNIX_EPOCH};
23
24use crate::compaction::leveled::{KeyRange, LeveledCompactionConfig, SSTableMeta};
25use crate::error::{Error, Result};
26use crate::kv::{KVStore, KVTransaction};
27use crate::lsm::buffer_pool::{BufferPool, BufferPoolConfig};
28use crate::lsm::checkpoint::{load_checkpoint_meta, save_checkpoint_meta, CheckpointMeta};
29use crate::lsm::memtable::{ImmutableMemTable, MemTable, MemTableConfig};
30use crate::lsm::metrics::{LsmMetrics, LsmMetricsSnapshot};
31use crate::lsm::sstable::{SSTableConfig, SSTableEntry, SSTableReader, SSTableWriter};
32use crate::lsm::wal::{
33 detect_wal_format_version, SyncMode, WalBatchOp, WalConfig, WalEntry, WalEntryPayload,
34 WalOpType, WalReader, WalWriter, WAL_FORMAT_VERSION,
35};
36use crate::storage::format::WriteThrottleConfig;
37use crate::txn::TxnManager;
38use crate::types::{Key, TxnId, TxnMode, Value};
39use tracing::{info, warn};
40
41#[derive(Debug, Clone, Copy, PartialEq, Eq)]
43pub enum ThreadMode {
44 MultiThread,
46 SingleThread,
48}
49
50#[derive(Debug, Clone)]
52pub struct LsmKVConfig {
53 pub wal: WalConfig,
55 pub checkpoint: CheckpointConfig,
57 pub memtable: MemTableConfig,
59 pub sstable: SSTableConfig,
61 pub compaction: LeveledCompactionConfig,
63 pub buffer_pool: BufferPoolConfig,
65 pub thread_mode: ThreadMode,
67 pub write_throttle: WriteThrottleConfig,
69}
70
71#[derive(Debug, Clone)]
73pub struct CheckpointConfig {
74 pub wal_size_threshold: u64,
76 pub min_interval_ms: u64,
78 pub auto_checkpoint: bool,
80}
81
82impl Default for CheckpointConfig {
83 fn default() -> Self {
84 Self {
85 wal_size_threshold: 64 * 1024 * 1024,
86 min_interval_ms: 60_000,
87 auto_checkpoint: true,
88 }
89 }
90}
91
92impl Default for LsmKVConfig {
93 fn default() -> Self {
94 let wal = WalConfig {
95 sync_mode: SyncMode::BatchSync {
96 max_batch_size: 1024,
97 max_wait_ms: 10,
98 },
99 ..WalConfig::default()
100 };
101
102 #[cfg(feature = "compression-lz4")]
104 let sstable = SSTableConfig {
105 compression: crate::lsm::sstable::CompressionType::Lz4,
106 ..SSTableConfig::default()
107 };
108 #[cfg(not(feature = "compression-lz4"))]
109 let sstable = SSTableConfig::default();
110
111 Self {
112 wal,
113 checkpoint: CheckpointConfig::default(),
114 memtable: MemTableConfig::default(),
115 sstable,
116 compaction: LeveledCompactionConfig::default(),
117 buffer_pool: BufferPoolConfig::default(),
118 thread_mode: ThreadMode::MultiThread,
119 write_throttle: WriteThrottleConfig::default(),
120 }
121 }
122}
123
124#[derive(Debug, Clone)]
126pub struct RecoveryResult {
127 pub entries_recovered: usize,
129 pub last_lsn: u64,
131 pub warnings: Vec<String>,
133 pub stop_reason: Option<String>,
135 pub checkpoint_lsn: Option<u64>,
137}
138
139#[derive(Debug, Clone)]
141pub struct CheckpointResult {
142 pub checkpoint_lsn: u64,
144 pub wal_bytes_reclaimed: u64,
146 pub duration_ms: u64,
148}
149
150#[derive(Debug)]
152pub struct TimestampOracle {
153 next: AtomicU64,
154}
155
156impl TimestampOracle {
157 pub fn new(start: u64) -> Self {
159 Self {
160 next: AtomicU64::new(start),
161 }
162 }
163
164 pub fn next_timestamp(&self) -> u64 {
166 self.next.fetch_add(1, Ordering::Relaxed)
167 }
168
169 pub fn current_timestamp(&self) -> u64 {
171 self.next.load(Ordering::Relaxed).saturating_sub(1)
172 }
173}
174
175#[derive(Debug)]
177pub struct LsmTxnManager {
178 next_txn_id: AtomicU64,
179}
180
181impl Default for LsmTxnManager {
182 fn default() -> Self {
183 Self {
184 next_txn_id: AtomicU64::new(1),
185 }
186 }
187}
188
189#[derive(Debug, Clone, Copy)]
190pub struct LsmTxnManagerRef<'a> {
192 store: &'a LsmKV,
193}
194
195impl<'a> LsmTxnManagerRef<'a> {
196 fn allocate_txn_id(&self) -> TxnId {
197 TxnId(
198 self.store
199 .txn_manager
200 .next_txn_id
201 .fetch_add(1, Ordering::Relaxed),
202 )
203 }
204}
205
206#[derive(Debug)]
210pub struct LsmKV {
211 pub config: LsmKVConfig,
213 pub data_dir: PathBuf,
215 pub sst_dir: PathBuf,
217 pub wal_path: PathBuf,
219 pub wal: RwLock<WalWriter>,
221 pub active_memtable: RwLock<MemTable>,
223 pub immutable_memtables: RwLock<VecDeque<Arc<ImmutableMemTable>>>,
225 pub levels: RwLock<Vec<Vec<SSTableMeta>>>,
227 pub buffer_pool: BufferPool,
229 pub metrics: Arc<LsmMetrics>,
231 pub ts_oracle: TimestampOracle,
233 pub txn_manager: LsmTxnManager,
235 pub commit_lock: Mutex<()>,
237 pub next_sstable_id: AtomicU64,
239 pub wal_used_bytes: AtomicU64,
241 pub last_checkpoint_ms: AtomicU64,
243}
244
245impl LsmKV {
246 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
250 let (store, _recovery) = Self::open_with_config(path, LsmKVConfig::default())?;
251 Ok(store)
252 }
253
254 pub fn open_with_config(
258 path: impl AsRef<Path>,
259 config: LsmKVConfig,
260 ) -> Result<(Self, RecoveryResult)> {
261 let data_dir = path.as_ref().to_path_buf();
262 fs::create_dir_all(&data_dir)?;
263 let wal_path = data_dir.join("lsm.wal");
264 let sst_dir = data_dir.join("sst");
265 fs::create_dir_all(&sst_dir)?;
266 let metrics = Arc::new(LsmMetrics::default());
267 let checkpoint_path = data_dir.join("checkpoint.meta");
268
269 let (wal_writer, recovered, next_ts, recovery, last_checkpoint_ms) = if wal_path.exists() {
270 let wal_version = detect_wal_format_version(&wal_path, &config.wal)?;
271 if wal_version < WAL_FORMAT_VERSION {
272 Self::migrate_legacy_wal(&wal_path, &checkpoint_path, &config.wal)?;
273 }
274 let start = Instant::now();
275 let checkpoint = load_checkpoint_meta(&checkpoint_path)?;
276 let checkpoint_lsn = checkpoint.as_ref().map(|meta| meta.checkpoint_lsn);
277 let last_checkpoint_ms = checkpoint.as_ref().map(|meta| meta.created_at).unwrap_or(0);
278 let mut reader = WalReader::open(&wal_path, config.wal.clone())?;
279 let replay = reader.replay()?;
280 let mut mem = MemTable::new();
281 let entries: Vec<_> = if let Some(start_lsn) = checkpoint_lsn {
282 replay
283 .entries
284 .into_iter()
285 .filter(|entry| entry.lsn > start_lsn)
286 .collect()
287 } else {
288 replay.entries
289 };
290 let mut last_lsn = apply_wal_replay(&mut mem, &entries);
291 if let Some(start_lsn) = checkpoint_lsn {
292 last_lsn = last_lsn.max(start_lsn);
293 }
294 let next = last_lsn.saturating_add(1).max(1);
295 for warning in &replay.warnings {
296 warn!(warning = %warning, "WAL recovery warning");
297 }
298 if let Some(reason) = &replay.stop_reason {
299 warn!(stop_reason = %reason, "WAL recovery stopped early");
300 }
301
302 let stopped_at = replay.stopped_at;
303 let recovery = RecoveryResult {
304 entries_recovered: entries.len(),
305 last_lsn,
306 warnings: replay.warnings,
307 stop_reason: replay.stop_reason,
308 checkpoint_lsn,
309 };
310 let duration_ms = start.elapsed().as_millis() as u64;
311 info!(
312 entries_recovered = recovery.entries_recovered,
313 checkpoint_lsn = ?recovery.checkpoint_lsn,
314 duration_ms,
315 "WAL recovery completed"
316 );
317 let mut wal_writer = WalWriter::open(&wal_path, config.wal.clone())?;
318 if let Some(valid_end) = stopped_at {
319 wal_writer.truncate_tail_to(valid_end)?;
320 }
321 (wal_writer, mem, next, recovery, last_checkpoint_ms)
322 } else {
323 (
324 WalWriter::create(&wal_path, config.wal.clone(), 1, 1)?,
325 MemTable::new(),
326 1,
327 RecoveryResult {
328 entries_recovered: 0,
329 last_lsn: 0,
330 warnings: Vec::new(),
331 stop_reason: None,
332 checkpoint_lsn: None,
333 },
334 0,
335 )
336 };
337 let wal_used_bytes = wal_writer.used_bytes();
338 let next_sstable_id = next_sstable_id_from_dir(&sst_dir)?;
339 let levels = load_sstable_levels(&sst_dir, config.compaction.max_levels)?;
340
341 let store = Self {
342 wal: RwLock::new(wal_writer),
343 active_memtable: RwLock::new(recovered),
344 immutable_memtables: RwLock::new(VecDeque::new()),
345 levels: RwLock::new(levels),
346 buffer_pool: BufferPool::new(config.buffer_pool),
347 metrics: Arc::clone(&metrics),
348 ts_oracle: TimestampOracle::new(next_ts),
349 txn_manager: LsmTxnManager::default(),
350 commit_lock: Mutex::new(()),
351 next_sstable_id: AtomicU64::new(next_sstable_id),
352 wal_used_bytes: AtomicU64::new(wal_used_bytes),
353 last_checkpoint_ms: AtomicU64::new(last_checkpoint_ms),
354 data_dir,
355 sst_dir,
356 wal_path,
357 config,
358 };
359 store.refresh_memtable_size_metrics();
360 Ok((store, recovery))
361 }
362
363 fn migrate_legacy_wal(
364 wal_path: &Path,
365 checkpoint_path: &Path,
366 config: &WalConfig,
367 ) -> Result<()> {
368 let mut reader = WalReader::open_allow_legacy(wal_path, config.clone())?;
369 let replay = reader.replay()?;
370 for warning in &replay.warnings {
371 warn!(warning = %warning, "Legacy WAL replay warning during migration");
372 }
373 if let Some(reason) = &replay.stop_reason {
374 warn!(
375 stop_reason = %reason,
376 "Legacy WAL replay stopped early during migration"
377 );
378 }
379
380 let entries = replay.entries;
381 let first_lsn = entries.first().map(|entry| entry.lsn).unwrap_or(1);
382 let temp_path = wal_path.with_extension("wal.migrate");
383 let backup_path = wal_path.with_extension("wal.bak");
384
385 if temp_path.exists() {
386 fs::remove_file(&temp_path)?;
387 }
388
389 let mut writer = WalWriter::create(&temp_path, config.clone(), 1, first_lsn)?;
390 for entry in &entries {
391 writer.append(entry)?;
392 }
393 writer.force_sync()?;
394 drop(writer);
395
396 if backup_path.exists() {
397 fs::remove_file(&backup_path)?;
398 }
399 fs::rename(wal_path, &backup_path)?;
400 fs::rename(&temp_path, wal_path)?;
401
402 let created_at = SystemTime::now()
403 .duration_since(UNIX_EPOCH)
404 .unwrap_or_default()
405 .as_millis() as u64;
406 let meta = CheckpointMeta::new(0, created_at);
407 save_checkpoint_meta(checkpoint_path, &meta)?;
408
409 info!(
410 entries_migrated = entries.len(),
411 backup = ?backup_path,
412 "Legacy WAL migration completed"
413 );
414
415 Ok(())
416 }
417
418 pub fn flush(&self) -> Result<()> {
422 let old = {
423 let mut guard = self
424 .active_memtable
425 .write()
426 .expect("lsm active_memtable lock poisoned");
427 std::mem::take(&mut *guard)
428 };
429 let imm = Arc::new(old.freeze());
430
431 {
432 let mut queue = self
433 .immutable_memtables
434 .write()
435 .expect("lsm immutable_memtables lock poisoned");
436 queue.push_back(imm);
437
438 while queue.len() > self.config.memtable.max_immutable_count {
439 queue.pop_front();
440 }
441 }
442 self.metrics.inc_memtable_flush_count();
443 self.refresh_memtable_size_metrics();
444 Ok(())
445 }
446
447 pub fn checkpoint(&self) -> Result<CheckpointResult> {
449 let _guard = self.commit_lock.lock().expect("lsm commit_lock poisoned");
450 let start = Instant::now();
451
452 self.flush()?;
453 self.persist_immutable_memtables()?;
454
455 let checkpoint_lsn = self.ts_oracle.current_timestamp();
456 let created_at = SystemTime::now()
457 .duration_since(UNIX_EPOCH)
458 .unwrap_or_default()
459 .as_millis() as u64;
460 let meta = CheckpointMeta::new(checkpoint_lsn, created_at);
461 let checkpoint_path = self.data_dir.join("checkpoint.meta");
462 save_checkpoint_meta(&checkpoint_path, &meta)?;
463
464 let mut wal = self.wal.write().expect("lsm wal lock poisoned");
465 let wal_bytes_reclaimed = wal.used_bytes();
466 let end_offset = wal.end_offset();
467 wal.advance_start(end_offset)?;
468 self.wal_used_bytes
469 .store(wal.used_bytes(), Ordering::Relaxed);
470 self.last_checkpoint_ms.store(created_at, Ordering::Relaxed);
471
472 Ok(CheckpointResult {
473 checkpoint_lsn,
474 wal_bytes_reclaimed,
475 duration_ms: start.elapsed().as_millis() as u64,
476 })
477 }
478
479 pub fn should_checkpoint(&self) -> bool {
481 if !self.config.checkpoint.auto_checkpoint {
482 return false;
483 }
484 let wal_used = self.wal_used_bytes.load(Ordering::Relaxed);
485 if wal_used <= self.config.checkpoint.wal_size_threshold {
486 return false;
487 }
488 let now_ms = SystemTime::now()
489 .duration_since(UNIX_EPOCH)
490 .unwrap_or_default()
491 .as_millis() as u64;
492 let last = self.last_checkpoint_ms.load(Ordering::Relaxed);
493 now_ms.saturating_sub(last) >= self.config.checkpoint.min_interval_ms
494 }
495
496 pub fn compact(&self) -> Result<()> {
500 Ok(())
501 }
502
503 pub fn metrics(&self) -> LsmMetricsSnapshot {
505 let counters = self.metrics.counters_snapshot();
507 let sstable_count_per_level = self
508 .levels
509 .read()
510 .expect("lsm levels lock poisoned")
511 .iter()
512 .map(|l| l.len())
513 .collect::<Vec<_>>();
514
515 let bp_stats = self.buffer_pool.stats();
516 let bp_total = bp_stats.hits + bp_stats.misses;
517 let hit_rate = if bp_total == 0 {
518 1.0
519 } else {
520 (bp_stats.hits as f64) / (bp_total as f64)
521 };
522 LsmMetricsSnapshot {
523 wal_write_bytes: counters.wal_write_bytes,
524 wal_sync_duration_ms: counters.wal_sync_duration_ms,
525 memtable_size_bytes: counters.memtable_size_bytes,
526 memtable_flush_count: counters.memtable_flush_count,
527 sstable_read_bytes: counters.sstable_read_bytes,
528 sstable_count_per_level,
529 buffer_pool_hit_rate: hit_rate,
530 buffer_pool_size_bytes: self.buffer_pool.current_size_bytes() as u64,
531 compaction_bytes_written: counters.compaction_bytes_written,
532 compaction_duration_ms: counters.compaction_duration_ms,
533 }
534 }
535
536 pub fn disk_usage(&self) -> u64 {
538 fs::metadata(&self.wal_path).map(|m| m.len()).unwrap_or(0)
539 }
540
541 fn sstable_path_for(&self, file_id: u64) -> PathBuf {
542 self.sst_dir.join(format!("{file_id}.sst"))
543 }
544
545 fn persist_immutable_memtables(&self) -> Result<()> {
546 let immutables = {
547 let mut guard = self
548 .immutable_memtables
549 .write()
550 .expect("lsm immutable_memtables lock poisoned");
551 std::mem::take(&mut *guard)
552 };
553
554 if immutables.is_empty() {
555 return Ok(());
556 }
557
558 let mut levels = self.levels.write().expect("lsm levels lock poisoned");
559 for mem in immutables {
560 let entries = mem.scan_prefix(b"", u64::MAX);
561 if entries.is_empty() {
562 continue;
563 }
564 let file_id = self.next_sstable_id.fetch_add(1, Ordering::Relaxed);
565 let path = self.sstable_path_for(file_id);
566 let mut writer = SSTableWriter::create(&path, self.config.sstable)?;
567
568 let mut first_key: Option<Key> = None;
569 let mut last_key: Option<Key> = None;
570 for (key, entry) in entries {
571 if first_key.is_none() {
572 first_key = Some(key.clone());
573 }
574 last_key = Some(key.clone());
575 writer.append(SSTableEntry {
576 key,
577 value: entry.value,
578 timestamp: entry.timestamp,
579 sequence: entry.sequence,
580 })?;
581 }
582 writer.finish()?;
583 let size_bytes = fs::metadata(&path)?.len();
584
585 let key_range = KeyRange {
586 first_key: first_key.unwrap(),
587 last_key: last_key.unwrap(),
588 };
589 let meta = SSTableMeta {
590 id: file_id,
591 level: 0,
592 size_bytes,
593 key_range,
594 };
595 levels[0].push(meta);
596 }
597
598 self.refresh_memtable_size_metrics();
599 Ok(())
600 }
601
602 fn refresh_memtable_size_metrics(&self) {
603 let active_bytes = self
604 .active_memtable
605 .read()
606 .expect("lsm active_memtable lock poisoned")
607 .memory_usage_bytes() as u64;
608 let imm_bytes = self
609 .immutable_memtables
610 .read()
611 .expect("lsm immutable_memtables lock poisoned")
612 .iter()
613 .map(|t| t.memory_usage_bytes() as u64)
614 .sum::<u64>();
615 self.metrics
616 .set_memtable_size_bytes(active_bytes.saturating_add(imm_bytes));
617 }
618
619 fn get_visible_at(
620 &self,
621 key: &[u8],
622 read_timestamp: u64,
623 ) -> Option<crate::lsm::memtable::MemTableEntry> {
624 if let Some(e) = self
625 .active_memtable
626 .read()
627 .expect("lsm active_memtable lock poisoned")
628 .get(key, read_timestamp)
629 {
630 return Some(e);
631 }
632 let imm = self
633 .immutable_memtables
634 .read()
635 .expect("lsm immutable_memtables lock poisoned");
636 for t in imm.iter().rev() {
637 if let Some(e) = t.get(key, read_timestamp) {
638 return Some(e);
639 }
640 }
641
642 let levels = self.levels.read().expect("lsm levels lock poisoned");
644 let mut best: Option<crate::lsm::memtable::MemTableEntry> = None;
645 for level in levels.iter() {
646 for meta in level.iter() {
647 let path = self.sstable_path_for(meta.id);
648 let Ok(mut reader) = SSTableReader::open(&path) else {
649 continue;
650 };
651 let Ok(found) = reader.get_with_buffer_pool(
652 &self.buffer_pool,
653 &self.metrics,
654 meta.id,
655 key,
656 read_timestamp,
657 ) else {
658 continue;
659 };
660 let Some(found) = found else {
661 continue;
662 };
663 let candidate = crate::lsm::memtable::MemTableEntry {
664 value: found.value,
665 timestamp: found.timestamp,
666 sequence: found.sequence,
667 };
668 let better = match &best {
669 None => true,
670 Some(cur) => {
671 (candidate.timestamp > cur.timestamp)
672 || (candidate.timestamp == cur.timestamp
673 && candidate.sequence > cur.sequence)
674 }
675 };
676 if better {
677 best = Some(candidate);
678 }
679 }
680 }
681 best
682 }
683
684 fn latest_timestamp(&self, key: &[u8]) -> u64 {
685 let mut best: Option<(u64, u64)> = None;
686 if let Some(e) = self.get_visible_at(key, u64::MAX) {
687 best = Some((e.timestamp, e.sequence));
688 }
689 match best {
690 Some((ts, _seq)) => ts,
691 None => 0,
692 }
693 }
694
695 fn scan_prefix_visible(
696 &self,
697 prefix: &[u8],
698 read_timestamp: u64,
699 ) -> BTreeMap<Key, crate::lsm::memtable::MemTableEntry> {
700 let mut out: BTreeMap<Key, crate::lsm::memtable::MemTableEntry> = BTreeMap::new();
701
702 let active = self
703 .active_memtable
704 .read()
705 .expect("lsm active_memtable lock poisoned");
706 for (k, e) in active.scan_prefix(prefix, read_timestamp) {
707 out.insert(k, e);
708 }
709
710 let imm = self
711 .immutable_memtables
712 .read()
713 .expect("lsm immutable_memtables lock poisoned");
714 for t in imm.iter().rev() {
715 for (k, e) in t.scan_prefix(prefix, read_timestamp) {
716 match out.get(&k) {
717 None => {
718 out.insert(k, e);
719 }
720 Some(cur) => {
721 let better = (e.timestamp > cur.timestamp)
722 || (e.timestamp == cur.timestamp && e.sequence > cur.sequence);
723 if better {
724 out.insert(k, e);
725 }
726 }
727 }
728 }
729 }
730
731 let levels = self.levels.read().expect("lsm levels lock poisoned");
733 for level in levels.iter() {
734 for meta in level.iter() {
735 let path = self.sstable_path_for(meta.id);
736 let Ok(mut reader) = SSTableReader::open(&path) else {
737 continue;
738 };
739 let Ok(entries) = reader.scan_prefix_with_buffer_pool(
740 &self.buffer_pool,
741 &self.metrics,
742 meta.id,
743 prefix,
744 read_timestamp,
745 ) else {
746 continue;
747 };
748 for e in entries {
749 let k = e.key.clone();
750 let candidate = crate::lsm::memtable::MemTableEntry {
751 value: e.value,
752 timestamp: e.timestamp,
753 sequence: e.sequence,
754 };
755 match out.get(&k) {
756 None => {
757 out.insert(k, candidate);
758 }
759 Some(cur) => {
760 let better = (candidate.timestamp > cur.timestamp)
761 || (candidate.timestamp == cur.timestamp
762 && candidate.sequence > cur.sequence);
763 if better {
764 out.insert(k, candidate);
765 }
766 }
767 }
768 }
769 }
770 }
771
772 out
773 }
774
775 fn scan_range_visible(
776 &self,
777 start: &[u8],
778 end: &[u8],
779 read_timestamp: u64,
780 ) -> BTreeMap<Key, crate::lsm::memtable::MemTableEntry> {
781 let mut out: BTreeMap<Key, crate::lsm::memtable::MemTableEntry> = BTreeMap::new();
782
783 let active = self
784 .active_memtable
785 .read()
786 .expect("lsm active_memtable lock poisoned");
787 for (k, e) in active.scan_range(start, end, read_timestamp) {
788 out.insert(k, e);
789 }
790
791 let imm = self
792 .immutable_memtables
793 .read()
794 .expect("lsm immutable_memtables lock poisoned");
795 for t in imm.iter().rev() {
796 for (k, e) in t.scan_range(start, end, read_timestamp) {
797 match out.get(&k) {
798 None => {
799 out.insert(k, e);
800 }
801 Some(cur) => {
802 let better = (e.timestamp > cur.timestamp)
803 || (e.timestamp == cur.timestamp && e.sequence > cur.sequence);
804 if better {
805 out.insert(k, e);
806 }
807 }
808 }
809 }
810 }
811
812 let levels = self.levels.read().expect("lsm levels lock poisoned");
813 for level in levels.iter() {
814 for meta in level.iter() {
815 let path = self.sstable_path_for(meta.id);
816 let Ok(mut reader) = SSTableReader::open(&path) else {
817 continue;
818 };
819 let Ok(entries) = reader.scan_range_with_buffer_pool(
820 &self.buffer_pool,
821 &self.metrics,
822 meta.id,
823 start,
824 end,
825 read_timestamp,
826 ) else {
827 continue;
828 };
829 for e in entries {
830 let k = e.key.clone();
831 let candidate = crate::lsm::memtable::MemTableEntry {
832 value: e.value,
833 timestamp: e.timestamp,
834 sequence: e.sequence,
835 };
836 match out.get(&k) {
837 None => {
838 out.insert(k, candidate);
839 }
840 Some(cur) => {
841 let better = (candidate.timestamp > cur.timestamp)
842 || (candidate.timestamp == cur.timestamp
843 && candidate.sequence > cur.sequence);
844 if better {
845 out.insert(k, candidate);
846 }
847 }
848 }
849 }
850 }
851 }
852
853 out
854 }
855}
856
857fn next_sstable_id_from_dir(dir: &Path) -> Result<u64> {
858 let mut max_id = 0u64;
859 if dir.exists() {
860 for entry in fs::read_dir(dir)? {
861 let entry = entry?;
862 let path = entry.path();
863 if path.extension().and_then(|s| s.to_str()) != Some("sst") {
864 continue;
865 }
866 let Some(stem) = path.file_stem().and_then(|s| s.to_str()) else {
867 continue;
868 };
869 let Ok(id) = stem.parse::<u64>() else {
870 continue;
871 };
872 max_id = max_id.max(id);
873 }
874 }
875 Ok(max_id.saturating_add(1))
876}
877
878fn load_sstable_levels(dir: &Path, max_levels: usize) -> Result<Vec<Vec<SSTableMeta>>> {
879 let mut levels = vec![Vec::new(); max_levels];
880 if max_levels == 0 || !dir.exists() {
881 return Ok(levels);
882 }
883
884 for entry in fs::read_dir(dir)? {
885 let entry = entry?;
886 let path = entry.path();
887 if path.extension().and_then(|s| s.to_str()) != Some("sst") {
888 continue;
889 }
890 let Some(stem) = path.file_stem().and_then(|s| s.to_str()) else {
891 continue;
892 };
893 let Ok(file_id) = stem.parse::<u64>() else {
894 continue;
895 };
896 let reader = match SSTableReader::open(&path) {
897 Ok(reader) => reader,
898 Err(err) => {
899 warn!(error = %err, path = ?path, "Skipping unreadable SSTable");
900 continue;
901 }
902 };
903 let Some((first_key, last_key)) = reader.key_range() else {
904 continue;
905 };
906 let size_bytes = fs::metadata(&path)?.len();
907 let meta = SSTableMeta {
908 id: file_id,
909 level: 0,
910 size_bytes,
911 key_range: KeyRange {
912 first_key,
913 last_key,
914 },
915 };
916 levels[0].push(meta);
917 }
918
919 Ok(levels)
920}
921
922fn apply_wal_replay(mem: &mut MemTable, entries: &[crate::lsm::wal::WalEntry]) -> u64 {
923 let mut last = 0u64;
924 for e in entries {
925 last = last.max(e.lsn);
926 match &e.payload {
927 WalEntryPayload::Put { key, value } => {
928 mem.put(key.clone(), value.clone(), e.lsn, 0);
929 }
930 WalEntryPayload::Delete { key } => {
931 mem.delete(key.clone(), e.lsn, 0);
932 }
933 WalEntryPayload::Batch(ops) => {
934 let mut seq = 0u64;
935 for op in ops {
936 match op.op_type {
937 crate::lsm::wal::WalOpType::Put => {
938 let val = op.value.clone().unwrap_or_default();
939 mem.put(op.key.clone(), val, e.lsn, seq);
940 }
941 crate::lsm::wal::WalOpType::Delete => {
942 mem.delete(op.key.clone(), e.lsn, seq);
943 }
944 }
945 seq = seq.wrapping_add(1);
946 }
947 }
948 }
949 }
950 last
951}
952
953#[derive(Debug)]
959pub struct LsmTransaction<'a> {
960 start_ts: u64,
962 txn_id: TxnId,
964 mode: TxnMode,
966 read_set: HashSet<Vec<u8>>,
968 write_set: BTreeMap<Vec<u8>, Option<Vec<u8>>>,
970 store: &'a LsmKV,
972}
973
974impl<'a> LsmTransaction<'a> {
975 fn new(store: &'a LsmKV, txn_id: TxnId, mode: TxnMode, start_ts: u64) -> Self {
976 Self {
977 start_ts,
978 txn_id,
979 mode,
980 read_set: HashSet::new(),
981 write_set: BTreeMap::new(),
982 store,
983 }
984 }
985
986 pub(crate) fn rollback_in_place(&mut self) -> Result<()> {
988 self.read_set.clear();
989 self.write_set.clear();
990 Ok(())
991 }
992
993 fn write_iter_prefix<'b>(
994 &'b self,
995 prefix: &'b [u8],
996 ) -> impl Iterator<Item = (&'b Key, &'b Option<Value>)> + 'b {
997 let prefix_vec = prefix.to_vec();
998 self.write_set
999 .range(prefix_vec..)
1000 .take_while(move |(k, _)| k.starts_with(prefix))
1001 }
1002}
1003
1004impl<'a> KVTransaction<'a> for LsmTransaction<'a> {
1005 fn id(&self) -> TxnId {
1006 self.txn_id
1007 }
1008
1009 fn mode(&self) -> TxnMode {
1010 self.mode
1011 }
1012
1013 fn get(&mut self, key: &Key) -> Result<Option<Value>> {
1014 if let Some(v) = self.write_set.get(key) {
1015 return Ok(v.clone());
1016 }
1017
1018 self.read_set.insert(key.clone());
1019 let entry = self.store.get_visible_at(key, self.start_ts);
1020 Ok(entry.and_then(|e| e.value))
1021 }
1022
1023 fn put(&mut self, key: Key, value: Value) -> Result<()> {
1024 if self.mode == TxnMode::ReadOnly {
1025 return Err(Error::TxnReadOnly);
1026 }
1027 self.write_set.insert(key, Some(value));
1028 Ok(())
1029 }
1030
1031 fn delete(&mut self, key: Key) -> Result<()> {
1032 if self.mode == TxnMode::ReadOnly {
1033 return Err(Error::TxnReadOnly);
1034 }
1035 self.write_set.insert(key, None);
1036 Ok(())
1037 }
1038
1039 fn scan_prefix(
1040 &mut self,
1041 prefix: &[u8],
1042 ) -> Result<Box<dyn Iterator<Item = (Key, Value)> + '_>> {
1043 let mut map: BTreeMap<Key, Option<Value>> = self
1044 .store
1045 .scan_prefix_visible(prefix, self.start_ts)
1046 .into_iter()
1047 .map(|(k, e)| (k, e.value))
1048 .collect();
1049
1050 self.read_set.extend(map.keys().cloned());
1052
1053 let overlays: Vec<(Key, Option<Value>)> = self
1054 .write_iter_prefix(prefix)
1055 .map(|(k, v)| (k.clone(), v.clone()))
1056 .collect();
1057 for (k, v) in overlays {
1058 self.read_set.insert(k.clone());
1059 match v {
1060 Some(val) => {
1061 map.insert(k, Some(val));
1062 }
1063 None => {
1064 map.remove(&k);
1065 }
1066 }
1067 }
1068
1069 let iter = map.into_iter().filter_map(|(k, v)| v.map(|vv| (k, vv)));
1070 Ok(Box::new(iter))
1071 }
1072
1073 fn scan_range(
1074 &mut self,
1075 start: &[u8],
1076 end: &[u8],
1077 ) -> Result<Box<dyn Iterator<Item = (Key, Value)> + '_>> {
1078 let mut map: BTreeMap<Key, Option<Value>> = self
1079 .store
1080 .scan_range_visible(start, end, self.start_ts)
1081 .into_iter()
1082 .map(|(k, e)| (k, e.value))
1083 .collect();
1084
1085 self.read_set.extend(map.keys().cloned());
1087
1088 let overlays: Vec<(Key, Option<Value>)> = self
1089 .write_set
1090 .range(start.to_vec()..end.to_vec())
1091 .map(|(k, v)| (k.clone(), v.clone()))
1092 .collect();
1093 for (k, v) in overlays {
1094 self.read_set.insert(k.clone());
1095 match v {
1096 Some(val) => {
1097 map.insert(k, Some(val));
1098 }
1099 None => {
1100 map.remove(&k);
1101 }
1102 }
1103 }
1104
1105 let iter = map.into_iter().filter_map(|(k, v)| v.map(|vv| (k, vv)));
1106 Ok(Box::new(iter))
1107 }
1108
1109 fn commit_self(mut self) -> Result<()> {
1110 if self.mode == TxnMode::ReadOnly || self.write_set.is_empty() {
1111 return Ok(());
1112 }
1113
1114 let _commit_guard = self
1115 .store
1116 .commit_lock
1117 .lock()
1118 .expect("lsm commit_lock poisoned");
1119
1120 for key in self.read_set.iter() {
1121 if self.store.latest_timestamp(key) > self.start_ts {
1122 return Err(Error::TxnConflict);
1123 }
1124 }
1125 for key in self.write_set.keys() {
1126 if self.store.latest_timestamp(key) > self.start_ts {
1127 return Err(Error::TxnConflict);
1128 }
1129 }
1130
1131 let commit_ts = self.store.ts_oracle.next_timestamp();
1132
1133 let mut ops = Vec::with_capacity(self.write_set.len());
1135 for (k, v) in self.write_set.iter() {
1136 let op = match v {
1137 Some(val) => WalBatchOp {
1138 op_type: WalOpType::Put,
1139 key: k.clone(),
1140 value: Some(val.clone()),
1141 },
1142 None => WalBatchOp {
1143 op_type: WalOpType::Delete,
1144 key: k.clone(),
1145 value: None,
1146 },
1147 };
1148 ops.push(op);
1149 }
1150 {
1151 let entry = WalEntry::batch(commit_ts, ops);
1152 let mut wal = self.store.wal.write().expect("lsm wal lock poisoned");
1153 let stats = wal.append_with_stats(&entry)?;
1154 self.store.metrics.add_wal_write_bytes(stats.bytes_written);
1155 let sync_duration_ms = if stats.sync_duration_ms == 0
1156 && !matches!(self.store.config.wal.sync_mode, SyncMode::EveryWrite)
1157 {
1158 wal.force_sync()?
1159 } else {
1160 stats.sync_duration_ms
1161 };
1162 self.store
1163 .metrics
1164 .add_wal_sync_duration_ms(sync_duration_ms);
1165 self.store
1166 .wal_used_bytes
1167 .store(wal.used_bytes(), Ordering::Relaxed);
1168 }
1169
1170 {
1171 let active = self
1172 .store
1173 .active_memtable
1174 .read()
1175 .expect("lsm active_memtable lock poisoned");
1176 let mut seq = 1u64;
1177 for (k, v) in std::mem::take(&mut self.write_set) {
1178 match v {
1179 Some(val) => active.put(k, val, commit_ts, seq),
1180 None => active.delete(k, commit_ts, seq),
1181 }
1182 seq = seq.wrapping_add(1);
1183 }
1184 }
1185 self.store.refresh_memtable_size_metrics();
1186
1187 if self
1189 .store
1190 .active_memtable
1191 .read()
1192 .expect("lsm active_memtable lock poisoned")
1193 .memory_usage_bytes()
1194 >= self.store.config.memtable.flush_threshold
1195 {
1196 self.store.flush()?;
1197 }
1198 Ok(())
1199 }
1200
1201 fn rollback_self(mut self) -> Result<()> {
1202 self.write_set.clear();
1203 Ok(())
1204 }
1205}
1206
1207impl<'a> TxnManager<'a, LsmTransaction<'a>> for LsmTxnManagerRef<'a> {
1208 fn begin(&'a self, mode: TxnMode) -> Result<LsmTransaction<'a>> {
1209 let start_ts = self.store.ts_oracle.next_timestamp();
1210 Ok(LsmTransaction::new(
1211 self.store,
1212 self.allocate_txn_id(),
1213 mode,
1214 start_ts,
1215 ))
1216 }
1217
1218 fn commit(&'a self, txn: LsmTransaction<'a>) -> Result<()> {
1219 txn.commit_self()
1220 }
1221
1222 fn rollback(&'a self, txn: LsmTransaction<'a>) -> Result<()> {
1223 txn.rollback_self()
1224 }
1225}
1226
1227impl KVStore for LsmKV {
1228 type Transaction<'a>
1229 = LsmTransaction<'a>
1230 where
1231 Self: 'a;
1232 type Manager<'a>
1233 = LsmTxnManagerRef<'a>
1234 where
1235 Self: 'a;
1236
1237 fn txn_manager(&self) -> Self::Manager<'_> {
1238 LsmTxnManagerRef { store: self }
1239 }
1240
1241 fn begin(&self, mode: TxnMode) -> Result<Self::Transaction<'_>> {
1242 let manager = LsmTxnManagerRef { store: self };
1243 let start_ts = self.ts_oracle.next_timestamp();
1244 Ok(LsmTransaction::new(
1245 self,
1246 manager.allocate_txn_id(),
1247 mode,
1248 start_ts,
1249 ))
1250 }
1251}
1252
1253#[cfg(all(test, not(target_arch = "wasm32")))]
1254mod kv_store {
1255 use super::*;
1256
1257 fn test_config() -> LsmKVConfig {
1258 LsmKVConfig {
1259 wal: WalConfig {
1260 segment_size: 4096,
1261 max_segments: 2,
1262 sync_mode: SyncMode::NoSync,
1263 },
1264 ..Default::default()
1265 }
1266 }
1267
1268 fn new_test_store() -> LsmKV {
1269 let cfg = test_config();
1270 let data_dir = tempfile::tempdir().expect("tempdir").keep();
1271 let sst_dir = data_dir.join("sst");
1272 fs::create_dir_all(&sst_dir).expect("create sst dir");
1273 let wal_path = data_dir.join("lsm.wal");
1274 let wal = WalWriter::create(&wal_path, cfg.wal.clone(), 1, 1).expect("wal create");
1275
1276 let levels = vec![Vec::new(); cfg.compaction.max_levels];
1277 LsmKV {
1278 config: cfg,
1279 data_dir,
1280 sst_dir,
1281 wal_path,
1282 wal: RwLock::new(wal),
1283 active_memtable: RwLock::new(MemTable::new()),
1284 immutable_memtables: RwLock::new(VecDeque::new()),
1285 levels: RwLock::new(levels),
1286 buffer_pool: BufferPool::new(BufferPoolConfig::default()),
1287 metrics: Arc::new(LsmMetrics::default()),
1288 ts_oracle: TimestampOracle::new(1),
1289 txn_manager: LsmTxnManager::default(),
1290 commit_lock: Mutex::new(()),
1291 next_sstable_id: AtomicU64::new(1),
1292 wal_used_bytes: AtomicU64::new(0),
1293 last_checkpoint_ms: AtomicU64::new(0),
1294 }
1295 }
1296
1297 fn new_test_store_with_sync(sync_mode: SyncMode) -> LsmKV {
1298 let mut cfg = test_config();
1299 cfg.wal.sync_mode = sync_mode;
1300 let data_dir = tempfile::tempdir().expect("tempdir").keep();
1301 let sst_dir = data_dir.join("sst");
1302 fs::create_dir_all(&sst_dir).expect("create sst dir");
1303 let wal_path = data_dir.join("lsm.wal");
1304 let wal = WalWriter::create(&wal_path, cfg.wal.clone(), 1, 1).expect("wal create");
1305
1306 let levels = vec![Vec::new(); cfg.compaction.max_levels];
1307 LsmKV {
1308 config: cfg,
1309 data_dir,
1310 sst_dir,
1311 wal_path,
1312 wal: RwLock::new(wal),
1313 active_memtable: RwLock::new(MemTable::new()),
1314 immutable_memtables: RwLock::new(VecDeque::new()),
1315 levels: RwLock::new(levels),
1316 buffer_pool: BufferPool::new(BufferPoolConfig::default()),
1317 metrics: Arc::new(LsmMetrics::default()),
1318 ts_oracle: TimestampOracle::new(1),
1319 txn_manager: LsmTxnManager::default(),
1320 commit_lock: Mutex::new(()),
1321 next_sstable_id: AtomicU64::new(1),
1322 wal_used_bytes: AtomicU64::new(0),
1323 last_checkpoint_ms: AtomicU64::new(0),
1324 }
1325 }
1326
1327 #[test]
1328 fn commit_forces_fsync_across_sync_modes() {
1329 let modes = [
1330 SyncMode::EveryWrite,
1331 SyncMode::BatchSync {
1332 max_batch_size: 1024 * 1024,
1333 max_wait_ms: 60_000,
1334 },
1335 SyncMode::NoSync,
1336 ];
1337
1338 for mode in modes {
1339 let store = new_test_store_with_sync(mode);
1340 crate::lsm::wal::reset_sync_calls();
1341
1342 let mut tx = store.begin(TxnMode::ReadWrite).unwrap();
1343 tx.put(b"k".to_vec(), b"v".to_vec()).unwrap();
1344 tx.commit_self().unwrap();
1345
1346 assert!(
1347 crate::lsm::wal::sync_calls() >= 1,
1348 "expected fsync during commit for sync mode"
1349 );
1350 }
1351 }
1352
1353 #[test]
1354 fn commit_makes_writes_visible() {
1355 let store = new_test_store();
1356 let mut tx = store.begin(TxnMode::ReadWrite).unwrap();
1357 tx.put(b"k".to_vec(), b"v".to_vec()).unwrap();
1358 assert_eq!(tx.get(&b"k".to_vec()).unwrap(), Some(b"v".to_vec()));
1359 tx.commit_self().unwrap();
1360
1361 let mut ro = store.begin(TxnMode::ReadOnly).unwrap();
1362 assert_eq!(ro.get(&b"k".to_vec()).unwrap(), Some(b"v".to_vec()));
1363 }
1364
1365 #[test]
1366 fn rollback_discards_writes() {
1367 let store = new_test_store();
1368 let mut tx = store.begin(TxnMode::ReadWrite).unwrap();
1369 tx.put(b"k".to_vec(), b"v".to_vec()).unwrap();
1370 tx.rollback_self().unwrap();
1371
1372 let mut ro = store.begin(TxnMode::ReadOnly).unwrap();
1373 assert_eq!(ro.get(&b"k".to_vec()).unwrap(), None);
1374 }
1375
1376 #[test]
1377 fn read_only_rejects_writes() {
1378 let store = new_test_store();
1379 let mut tx = store.begin(TxnMode::ReadOnly).unwrap();
1380 assert!(tx.put(b"k".to_vec(), b"v".to_vec()).is_err());
1381 }
1382
1383 #[test]
1384 fn detects_write_conflict() {
1385 let store = new_test_store();
1386 let mut a = store.begin(TxnMode::ReadWrite).unwrap();
1387 let mut b = store.begin(TxnMode::ReadWrite).unwrap();
1388
1389 a.put(b"k".to_vec(), b"v1".to_vec()).unwrap();
1390 a.commit_self().unwrap();
1391
1392 b.put(b"k".to_vec(), b"v2".to_vec()).unwrap();
1393 assert!(b.commit_self().is_err());
1394 }
1395
1396 #[test]
1397 fn scan_populates_read_set_for_conflict_detection() {
1398 let store = new_test_store();
1399
1400 let mut init = store.begin(TxnMode::ReadWrite).unwrap();
1401 init.put(b"p:a".to_vec(), b"v1".to_vec()).unwrap();
1402 init.commit_self().unwrap();
1403
1404 let mut scan_tx = store.begin(TxnMode::ReadWrite).unwrap();
1405 let got: Vec<(Key, Value)> = scan_tx.scan_prefix(b"p:").unwrap().collect();
1406 assert_eq!(got.len(), 1);
1407 assert_eq!(got[0].0, b"p:a".to_vec());
1408 assert_eq!(got[0].1, b"v1".to_vec());
1409
1410 let mut updater = store.begin(TxnMode::ReadWrite).unwrap();
1411 updater.put(b"p:a".to_vec(), b"v2".to_vec()).unwrap();
1412 updater.commit_self().unwrap();
1413
1414 scan_tx.put(b"q:z".to_vec(), b"ok".to_vec()).unwrap();
1415 assert!(scan_tx.commit_self().is_err());
1416 }
1417}
1418
1419#[cfg(all(test, not(target_arch = "wasm32")))]
1420mod txn {
1421 use super::*;
1422
1423 fn test_config() -> LsmKVConfig {
1424 LsmKVConfig {
1425 wal: WalConfig {
1426 segment_size: 4096,
1427 max_segments: 2,
1428 sync_mode: SyncMode::NoSync,
1429 },
1430 ..Default::default()
1431 }
1432 }
1433
1434 fn new_test_store() -> LsmKV {
1435 let cfg = test_config();
1436 let data_dir = tempfile::tempdir().expect("tempdir").keep();
1437 let sst_dir = data_dir.join("sst");
1438 fs::create_dir_all(&sst_dir).expect("create sst dir");
1439 let wal_path = data_dir.join("lsm.wal");
1440 let wal = WalWriter::create(&wal_path, cfg.wal.clone(), 1, 1).expect("wal create");
1441
1442 let levels = vec![Vec::new(); cfg.compaction.max_levels];
1443 LsmKV {
1444 config: cfg,
1445 data_dir,
1446 sst_dir,
1447 wal_path,
1448 wal: RwLock::new(wal),
1449 active_memtable: RwLock::new(MemTable::new()),
1450 immutable_memtables: RwLock::new(VecDeque::new()),
1451 levels: RwLock::new(levels),
1452 buffer_pool: BufferPool::new(BufferPoolConfig::default()),
1453 metrics: Arc::new(LsmMetrics::default()),
1454 ts_oracle: TimestampOracle::new(1),
1455 txn_manager: LsmTxnManager::default(),
1456 commit_lock: Mutex::new(()),
1457 next_sstable_id: AtomicU64::new(1),
1458 wal_used_bytes: AtomicU64::new(0),
1459 last_checkpoint_ms: AtomicU64::new(0),
1460 }
1461 }
1462
1463 #[test]
1464 fn detects_read_write_conflict_after_get() {
1465 let store = new_test_store();
1466
1467 let mut init = store.begin(TxnMode::ReadWrite).unwrap();
1468 init.put(b"k".to_vec(), b"v1".to_vec()).unwrap();
1469 init.commit_self().unwrap();
1470
1471 let mut a = store.begin(TxnMode::ReadWrite).unwrap();
1472 assert_eq!(a.get(&b"k".to_vec()).unwrap(), Some(b"v1".to_vec()));
1473
1474 let mut b = store.begin(TxnMode::ReadWrite).unwrap();
1475 b.put(b"k".to_vec(), b"v2".to_vec()).unwrap();
1476 b.commit_self().unwrap();
1477
1478 a.put(b"other".to_vec(), b"x".to_vec()).unwrap();
1479 assert!(a.commit_self().is_err());
1480 }
1481
1482 #[test]
1483 fn rollback_discards_write_set() {
1484 let store = new_test_store();
1485 let mut tx = store.begin(TxnMode::ReadWrite).unwrap();
1486 tx.put(b"k".to_vec(), b"v".to_vec()).unwrap();
1487 tx.rollback_self().unwrap();
1488
1489 let mut ro = store.begin(TxnMode::ReadOnly).unwrap();
1490 assert_eq!(ro.get(&b"k".to_vec()).unwrap(), None);
1491 }
1492}
1493
1494#[cfg(all(test, not(target_arch = "wasm32")))]
1495mod methods {
1496 use super::*;
1497
1498 fn test_config() -> LsmKVConfig {
1499 LsmKVConfig {
1500 wal: WalConfig {
1501 segment_size: 4096,
1502 max_segments: 2,
1503 sync_mode: SyncMode::NoSync,
1504 },
1505 ..Default::default()
1506 }
1507 }
1508
1509 #[test]
1510 fn open_creates_wal_and_returns_metrics() {
1511 let dir = tempfile::tempdir().expect("tempdir");
1512 let (store, _recovery) = LsmKV::open_with_config(dir.path(), test_config()).expect("open");
1513 let m = store.metrics();
1514 assert_eq!(m.wal_write_bytes, 0);
1515 assert_eq!(m.memtable_flush_count, 0);
1516 assert!(store.disk_usage() > 0);
1517 }
1518
1519 #[test]
1520 fn open_replays_wal_entries() {
1521 let dir = tempfile::tempdir().expect("tempdir");
1522
1523 {
1524 let (store, _recovery) =
1525 LsmKV::open_with_config(dir.path(), test_config()).expect("open");
1526 let mut wal = store.wal.write().unwrap();
1527 wal.append(&crate::lsm::wal::WalEntry::put(
1528 10,
1529 b"k".to_vec(),
1530 b"v".to_vec(),
1531 ))
1532 .unwrap();
1533 }
1534
1535 let (store, _recovery) =
1536 LsmKV::open_with_config(dir.path(), test_config()).expect("reopen");
1537 let mut tx = store.begin(TxnMode::ReadOnly).unwrap();
1538 assert_eq!(tx.get(&b"k".to_vec()).unwrap(), Some(b"v".to_vec()));
1539 }
1540
1541 #[test]
1542 fn flush_moves_active_to_immutable() {
1543 let dir = tempfile::tempdir().expect("tempdir");
1544 let (store, _recovery) = LsmKV::open_with_config(dir.path(), test_config()).expect("open");
1545
1546 let mut tx = store.begin(TxnMode::ReadWrite).unwrap();
1547 tx.put(b"k".to_vec(), b"v".to_vec()).unwrap();
1548 tx.commit_self().unwrap();
1549
1550 store.flush().unwrap();
1551
1552 let m = store.metrics();
1553 assert_eq!(m.memtable_flush_count, 1);
1554
1555 let mut ro = store.begin(TxnMode::ReadOnly).unwrap();
1556 assert_eq!(ro.get(&b"k".to_vec()).unwrap(), Some(b"v".to_vec()));
1557 }
1558}
1559
1560#[cfg(all(test, not(target_arch = "wasm32")))]
1561mod write_path {
1562 use super::*;
1563
1564 fn test_config() -> LsmKVConfig {
1565 LsmKVConfig {
1566 wal: WalConfig {
1567 segment_size: 4096,
1568 max_segments: 2,
1569 sync_mode: SyncMode::NoSync,
1570 },
1571 memtable: MemTableConfig {
1572 flush_threshold: 1,
1573 ..Default::default()
1574 },
1575 ..Default::default()
1576 }
1577 }
1578
1579 #[test]
1580 fn commit_appends_wal_and_reopen_replays() {
1581 let dir = tempfile::tempdir().expect("tempdir");
1582 {
1583 let (store, _recovery) =
1584 LsmKV::open_with_config(dir.path(), test_config()).expect("open");
1585 let mut tx = store.begin(TxnMode::ReadWrite).unwrap();
1586 tx.put(b"k".to_vec(), b"v".to_vec()).unwrap();
1587 tx.commit_self().unwrap();
1588 }
1589
1590 let (store, _recovery) =
1591 LsmKV::open_with_config(dir.path(), test_config()).expect("reopen");
1592 let mut ro = store.begin(TxnMode::ReadOnly).unwrap();
1593 assert_eq!(ro.get(&b"k".to_vec()).unwrap(), Some(b"v".to_vec()));
1594 }
1595
1596 #[test]
1597 fn flush_trigger_moves_to_immutable() {
1598 let dir = tempfile::tempdir().expect("tempdir");
1599 let (store, _recovery) = LsmKV::open_with_config(dir.path(), test_config()).expect("open");
1600
1601 let mut tx = store.begin(TxnMode::ReadWrite).unwrap();
1602 tx.put(b"k".to_vec(), b"v".to_vec()).unwrap();
1603 tx.commit_self().unwrap();
1604
1605 let metrics = store.metrics();
1607 assert_eq!(metrics.memtable_flush_count, 1);
1608
1609 let mut ro = store.begin(TxnMode::ReadOnly).unwrap();
1610 assert_eq!(ro.get(&b"k".to_vec()).unwrap(), Some(b"v".to_vec()));
1611 }
1612}
1613
1614#[cfg(all(test, not(target_arch = "wasm32")))]
1615mod recovery_tests {
1616 use super::*;
1617 use crate::lsm::checkpoint::{save_checkpoint_meta, CheckpointMeta};
1618 use std::io::{Read, Seek, SeekFrom, Write};
1619
1620 fn test_config() -> LsmKVConfig {
1621 LsmKVConfig {
1622 wal: WalConfig {
1623 segment_size: 4096,
1624 max_segments: 2,
1625 sync_mode: SyncMode::NoSync,
1626 },
1627 ..Default::default()
1628 }
1629 }
1630
1631 fn create_corrupted_tail_wal(dir: &Path, wal_cfg: WalConfig) {
1632 let wal_path = dir.join("lsm.wal");
1633 let mut writer = WalWriter::create(&wal_path, wal_cfg, 1, 1).unwrap();
1634 let e1 = WalEntry::put(1, b"a".to_vec(), b"1".to_vec());
1635 let e2 = WalEntry::put(2, b"b".to_vec(), b"2".to_vec());
1636 let _off1 = writer.append(&e1).unwrap();
1637 let off2 = writer.append(&e2).unwrap();
1638 let e2_bytes = e2.encode().unwrap();
1639 drop(writer);
1640
1641 let mut file = std::fs::OpenOptions::new()
1642 .read(true)
1643 .write(true)
1644 .open(&wal_path)
1645 .unwrap();
1646 let corrupt_offset = off2 + (e2_bytes.len() as u64).saturating_sub(1);
1647 file.seek(SeekFrom::Start(corrupt_offset)).unwrap();
1648 let mut buf = [0u8; 1];
1649 file.read_exact(&mut buf).unwrap();
1650 buf[0] ^= 0xFF;
1651 file.seek(SeekFrom::Start(corrupt_offset)).unwrap();
1652 file.write_all(&buf).unwrap();
1653 file.flush().unwrap();
1654 }
1655
1656 #[test]
1657 fn recovery_uses_checkpoint_lsn_when_present() {
1658 let dir = tempfile::tempdir().expect("tempdir");
1659 let (store, _recovery) = LsmKV::open_with_config(dir.path(), test_config()).expect("open");
1660 let mut tx = store.begin(TxnMode::ReadWrite).unwrap();
1661 tx.put(b"before".to_vec(), b"1".to_vec()).unwrap();
1662 tx.commit_self().unwrap();
1663
1664 let checkpoint_lsn = store.ts_oracle.current_timestamp();
1665 let meta = CheckpointMeta::new(checkpoint_lsn, 0);
1666 let checkpoint_path = dir.path().join("checkpoint.meta");
1667 save_checkpoint_meta(&checkpoint_path, &meta).unwrap();
1668
1669 let mut tx = store.begin(TxnMode::ReadWrite).unwrap();
1670 tx.put(b"after".to_vec(), b"2".to_vec()).unwrap();
1671 tx.commit_self().unwrap();
1672
1673 let (store, recovery) = LsmKV::open_with_config(dir.path(), test_config()).expect("reopen");
1674 assert!(recovery.checkpoint_lsn.is_some());
1675 assert_eq!(recovery.entries_recovered, 1);
1676
1677 let mut ro = store.begin(TxnMode::ReadOnly).unwrap();
1678 assert_eq!(ro.get(&b"after".to_vec()).unwrap(), Some(b"2".to_vec()));
1679 }
1680
1681 #[test]
1682 fn recovery_falls_back_when_checkpoint_missing() {
1683 let dir = tempfile::tempdir().expect("tempdir");
1684 let (store, _recovery) = LsmKV::open_with_config(dir.path(), test_config()).expect("open");
1685 let mut tx = store.begin(TxnMode::ReadWrite).unwrap();
1686 tx.put(b"k".to_vec(), b"v".to_vec()).unwrap();
1687 tx.commit_self().unwrap();
1688
1689 let (store, recovery) = LsmKV::open_with_config(dir.path(), test_config()).expect("reopen");
1690 assert!(recovery.checkpoint_lsn.is_none());
1691 assert!(recovery.entries_recovered >= 1);
1692
1693 let mut ro = store.begin(TxnMode::ReadOnly).unwrap();
1694 assert_eq!(ro.get(&b"k".to_vec()).unwrap(), Some(b"v".to_vec()));
1695 }
1696
1697 #[test]
1698 fn recovery_stops_on_corrupted_entry() {
1699 let dir = tempfile::tempdir().expect("tempdir");
1700 let wal_cfg = WalConfig {
1701 segment_size: 4096,
1702 max_segments: 1,
1703 sync_mode: SyncMode::NoSync,
1704 };
1705 create_corrupted_tail_wal(dir.path(), wal_cfg.clone());
1706
1707 let cfg = LsmKVConfig {
1708 wal: wal_cfg,
1709 ..Default::default()
1710 };
1711 let (store, recovery) = LsmKV::open_with_config(dir.path(), cfg).expect("reopen");
1712 assert!(recovery.stop_reason.is_some());
1713 assert_eq!(recovery.entries_recovered, 1);
1714 assert_eq!(recovery.last_lsn, 1);
1715
1716 let mut ro = store.begin(TxnMode::ReadOnly).unwrap();
1717 assert_eq!(ro.get(&b"a".to_vec()).unwrap(), Some(b"1".to_vec()));
1718 }
1719
1720 #[test]
1721 fn recovery_is_idempotent_across_reopens() {
1722 let dir = tempfile::tempdir().expect("tempdir");
1723 let wal_cfg = WalConfig {
1724 segment_size: 4096,
1725 max_segments: 1,
1726 sync_mode: SyncMode::NoSync,
1727 };
1728 create_corrupted_tail_wal(dir.path(), wal_cfg.clone());
1729
1730 let cfg = LsmKVConfig {
1731 wal: wal_cfg,
1732 ..Default::default()
1733 };
1734
1735 let (first_recovery, first_data) = {
1736 let (store, recovery) =
1737 LsmKV::open_with_config(dir.path(), cfg.clone()).expect("first reopen");
1738 let mut ro = store.begin(TxnMode::ReadOnly).unwrap();
1739 let data = (
1740 ro.get(&b"a".to_vec()).unwrap(),
1741 ro.get(&b"b".to_vec()).unwrap(),
1742 );
1743 (recovery, data)
1744 };
1745
1746 let (second_recovery, second_data) = {
1747 let (store, recovery) =
1748 LsmKV::open_with_config(dir.path(), cfg).expect("second reopen");
1749 let mut ro = store.begin(TxnMode::ReadOnly).unwrap();
1750 let data = (
1751 ro.get(&b"a".to_vec()).unwrap(),
1752 ro.get(&b"b".to_vec()).unwrap(),
1753 );
1754 (recovery, data)
1755 };
1756
1757 assert_eq!(
1758 first_recovery.entries_recovered,
1759 second_recovery.entries_recovered
1760 );
1761 assert_eq!(first_recovery.last_lsn, second_recovery.last_lsn);
1762 assert_eq!(first_recovery.entries_recovered, 1);
1763 assert_eq!(first_recovery.last_lsn, 1);
1764 assert_eq!(first_data, second_data);
1765 assert_eq!(first_data, (Some(b"1".to_vec()), None));
1766 }
1767}
1768
1769#[cfg(all(test, not(target_arch = "wasm32")))]
1770mod read_path {
1771 use super::*;
1772
1773 use crate::compaction::leveled::KeyRange;
1774 use crate::lsm::sstable::{SSTableEntry, SSTableWriter};
1775
1776 fn test_config() -> LsmKVConfig {
1777 LsmKVConfig {
1778 wal: WalConfig {
1779 segment_size: 4096,
1780 max_segments: 2,
1781 sync_mode: SyncMode::NoSync,
1782 },
1783 ..Default::default()
1784 }
1785 }
1786
1787 #[test]
1788 fn reads_from_sstable_via_buffer_pool() {
1789 let dir = tempfile::tempdir().expect("tempdir");
1790 let (store, _recovery) = LsmKV::open_with_config(dir.path(), test_config()).expect("open");
1791
1792 let file_id = 1u64;
1794 let path = store.sstable_path_for(file_id);
1795 let mut writer = SSTableWriter::create(&path, store.config.sstable).expect("sst create");
1796 writer
1797 .append(SSTableEntry {
1798 key: b"k".to_vec(),
1799 value: Some(b"v".to_vec()),
1800 timestamp: 0,
1801 sequence: 1,
1802 })
1803 .unwrap();
1804 writer.finish().unwrap();
1805
1806 let size_bytes = fs::metadata(&path).unwrap().len();
1807 let meta = SSTableMeta {
1808 id: file_id,
1809 level: 0,
1810 size_bytes,
1811 key_range: KeyRange {
1812 first_key: b"k".to_vec(),
1813 last_key: b"k".to_vec(),
1814 },
1815 };
1816 store.levels.write().unwrap()[0].push(meta);
1817
1818 let before = store.buffer_pool.stats();
1819 let mut ro = store.begin(TxnMode::ReadOnly).unwrap();
1820 assert_eq!(ro.get(&b"k".to_vec()).unwrap(), Some(b"v".to_vec()));
1821 assert_eq!(ro.get(&b"k".to_vec()).unwrap(), Some(b"v".to_vec()));
1822 let after = store.buffer_pool.stats();
1823
1824 assert!(after.misses > before.misses);
1825 assert!(after.hits > before.hits);
1826 }
1827}
1828
1829#[cfg(all(test, not(target_arch = "wasm32")))]
1830mod integration;