Skip to main content

alopex_core/lsm/
mod.rs

1//! ディスク永続化向けの LSM-Tree 実装。
2//!
3//! このモジュールは「単一 `.alopex` ファイル」方針の Disk モード向けに、WAL / MemTable /
4//! SSTable / Compaction を統合する `LsmKV` の土台を提供する。
5//!
6//! 仕様: `docs-internal/specs/lsm-tree-file-mode-spec.md`
7
8pub 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/// スレッドアクセスモード。
42#[derive(Debug, Clone, Copy, PartialEq, Eq)]
43pub enum ThreadMode {
44    /// マルチスレッド同時アクセス(デフォルト)。
45    MultiThread,
46    /// シングルスレッド専有アクセス(ロックオーバーヘッド最小)。
47    SingleThread,
48}
49
50/// LSM-Tree の設定。
51#[derive(Debug, Clone)]
52pub struct LsmKVConfig {
53    /// WAL 設定。
54    pub wal: WalConfig,
55    /// チェックポイント設定。
56    pub checkpoint: CheckpointConfig,
57    /// MemTable 設定。
58    pub memtable: MemTableConfig,
59    /// SSTable 設定。
60    pub sstable: SSTableConfig,
61    /// Compaction 設定。
62    pub compaction: LeveledCompactionConfig,
63    /// バッファプール設定。
64    pub buffer_pool: BufferPoolConfig,
65    /// スレッドモード。
66    pub thread_mode: ThreadMode,
67    /// 書き込みスロットリング設定。
68    pub write_throttle: WriteThrottleConfig,
69}
70
71/// チェックポイント設定。
72#[derive(Debug, Clone)]
73pub struct CheckpointConfig {
74    /// WAL サイズ閾値(バイト)。
75    pub wal_size_threshold: u64,
76    /// 最小チェックポイント間隔(ms)。
77    pub min_interval_ms: u64,
78    /// 自動チェックポイントを有効にするか。
79    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        // 仕様書のデフォルトは LZ4 だが、機能フラグ未指定でもコンパイルできるように分岐する。
103        #[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/// WAL リカバリの診断結果。
125#[derive(Debug, Clone)]
126pub struct RecoveryResult {
127    /// リカバリで適用したエントリ数。
128    pub entries_recovered: usize,
129    /// 最後に適用した LSN。
130    pub last_lsn: u64,
131    /// 非致命的な警告メッセージ。
132    pub warnings: Vec<String>,
133    /// リカバリ停止理由(ある場合)。
134    pub stop_reason: Option<String>,
135    /// チェックポイント LSN(利用時のみ)。
136    pub checkpoint_lsn: Option<u64>,
137}
138
139/// チェックポイント実行結果。
140#[derive(Debug, Clone)]
141pub struct CheckpointResult {
142    /// Checkpoint LSN captured during the run.
143    pub checkpoint_lsn: u64,
144    /// WAL bytes reclaimed by advancing the start offset.
145    pub wal_bytes_reclaimed: u64,
146    /// Total checkpoint duration in milliseconds.
147    pub duration_ms: u64,
148}
149
150/// タイムスタンプ生成器(単調増加)。
151#[derive(Debug)]
152pub struct TimestampOracle {
153    next: AtomicU64,
154}
155
156impl TimestampOracle {
157    /// 新しいオラクルを作成する。
158    pub fn new(start: u64) -> Self {
159        Self {
160            next: AtomicU64::new(start),
161        }
162    }
163
164    /// 新しいタイムスタンプを発行する。
165    pub fn next_timestamp(&self) -> u64 {
166        self.next.fetch_add(1, Ordering::Relaxed)
167    }
168
169    /// Latest issued timestamp without incrementing.
170    pub fn current_timestamp(&self) -> u64 {
171        self.next.load(Ordering::Relaxed).saturating_sub(1)
172    }
173}
174
175/// LSM 用トランザクションマネージャ(詳細はタスク 3.3 で実装)。
176#[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)]
190/// `LsmKV` に紐づくトランザクションマネージャの参照。
191pub 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/// LSM-Tree ベースの KV ストア(Disk モード)。
207///
208/// 設計: 仕様書 §4.1
209#[derive(Debug)]
210pub struct LsmKV {
211    /// 設定。
212    pub config: LsmKVConfig,
213    /// データディレクトリ。
214    pub data_dir: PathBuf,
215    /// SSTable ディレクトリ。
216    pub sst_dir: PathBuf,
217    /// WAL ファイルパス。
218    pub wal_path: PathBuf,
219    /// WAL Writer。
220    pub wal: RwLock<WalWriter>,
221    /// アクティブ MemTable。
222    pub active_memtable: RwLock<MemTable>,
223    /// Immutable MemTable キュー。
224    pub immutable_memtables: RwLock<VecDeque<Arc<ImmutableMemTable>>>,
225    /// レベル別 SSTable 一覧(コンパクションの単位)。
226    pub levels: RwLock<Vec<Vec<SSTableMeta>>>,
227    /// SSTable データブロックのバッファプール。
228    pub buffer_pool: BufferPool,
229    /// メトリクス(Atomic カウンタ)。
230    pub metrics: Arc<LsmMetrics>,
231    /// タイムスタンプオラクル。
232    pub ts_oracle: TimestampOracle,
233    /// トランザクションマネージャ。
234    pub txn_manager: LsmTxnManager,
235    /// コミットの直列化ロック(OCC の検証ウィンドウを閉じる)。
236    pub commit_lock: Mutex<()>,
237    /// 次に割り当てる SSTable ID。
238    pub next_sstable_id: AtomicU64,
239    /// WAL の現在使用量(バイト)。
240    pub wal_used_bytes: AtomicU64,
241    /// 最終チェックポイント時刻(epoch ms)。
242    pub last_checkpoint_ms: AtomicU64,
243}
244
245impl LsmKV {
246    /// LsmKV をデフォルト設定で開く。
247    ///
248    /// `path` はデータディレクトリとして扱い、内部で WAL ファイルを作成/再利用する。
249    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    /// 設定付きで LsmKV を開く。
255    ///
256    /// 既存 WAL がある場合は WAL をリプレイして MemTable を復元する(クラッシュリカバリ)。
257    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    /// MemTable をフラッシュする(手動)。
419    ///
420    /// 現段階では「Active MemTable を freeze して Immutable キューへ移す」までを行う。
421    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    /// 明示的にチェックポイントを作成する。
448    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    /// Determine whether an auto-checkpoint should run based on size and time thresholds.
480    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    /// コンパクションを実行する(手動)。
497    ///
498    /// 現段階ではメタデータ更新までの配線は未実装のため、no-op とする。
499    pub fn compact(&self) -> Result<()> {
500        Ok(())
501    }
502
503    /// メトリクスを取得する。
504    pub fn metrics(&self) -> LsmMetricsSnapshot {
505        // Atomic カウンタの値を取得し、スナップショットに補完情報を足す。
506        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    /// 推定ディスク使用量(バイト)を返す。
537    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        // SSTable を探索(L0 → L1..Ln)。
643        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        // SSTable もマージ(L0→Ln、同一キーは新しい timestamp/sequence を採用)。
732        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/// LSM 用トランザクション(Snapshot Isolation + OCC)。
954///
955/// - 開始時点の `start_ts` をスナップショットとして読み取る。
956/// - `read_set` を記録し、コミット時に read-write conflict を検出する(ファントム検出は未対応)。
957/// - 書き込みは `write_set` にバッファし、コミット時に WAL → MemTable の順で反映する。
958#[derive(Debug)]
959pub struct LsmTransaction<'a> {
960    /// 開始タイムスタンプ。
961    start_ts: u64,
962    /// トランザクション ID。
963    txn_id: TxnId,
964    /// モード。
965    mode: TxnMode,
966    /// 読み取りセット。
967    read_set: HashSet<Vec<u8>>,
968    /// 書き込みセット(`None` は tombstone)。
969    write_set: BTreeMap<Vec<u8>, Option<Vec<u8>>>,
970    /// KV ストア参照。
971    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    /// トランザクションを消費せずにロールバックする。
987    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        // スナップショットで観測したキーは read_set に入れる(read-write conflict 検出の最低限)。
1051        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        // スナップショットで観測したキーは read_set に入れる(read-write conflict 検出の最低限)。
1086        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        // WAL に先行して追記してから MemTable に反映する(WAL → MemTable)。
1134        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        // flush trigger: MemTable が閾値を超えたら freeze して immutable へ移動。
1188        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        // flush_threshold=1 のため commit 後に自動で flush される。
1606        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        // SSTable を1つ作成して L0 に登録する。
1793        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;