Skip to main content

alopex_core/kv/
memory.rs

1//! An in-memory key-value store implementation with Write-Ahead Logging
2//! and Optimistic Concurrency Control for Snapshot Isolation.
3
4use crate::error::{Error, Result};
5#[cfg(feature = "test-hooks")]
6use crate::kv::hooks::{CrashOperation, CrashSimulator, CrashTiming, IoHooks};
7use crate::kv::{KVStore, KVTransaction, OwnedKVScan, OwnedKVStore, OwnedKVTransaction};
8use crate::log::wal::{WalReader, WalRecord, WalWriter};
9use crate::storage::flush::write_empty_vector_segment;
10use crate::storage::sstable::{SstableReader, SstableWriter};
11use crate::txn::TxnManager;
12use crate::types::{Key, TxnId, TxnMode, TxnState, Value};
13use std::collections::{BTreeMap, HashMap};
14use std::ops::Bound::{Excluded, Included, Unbounded};
15use std::path::{Path, PathBuf};
16use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
17use std::sync::{Arc, Condvar, Mutex, RwLock, RwLockReadGuard};
18use tracing::warn;
19
20/// メモリ使用量の統計(バイト単位)。
21#[derive(Debug, Clone, Default, PartialEq)]
22pub struct MemoryStats {
23    /// 全体のメモリ使用量。
24    pub total_bytes: usize,
25    /// KV データのメモリ使用量。
26    pub kv_bytes: usize,
27    /// 補助インデックスのメモリ使用量。
28    pub index_bytes: usize,
29}
30
31/// An in-memory key-value store.
32#[derive(Clone)]
33pub struct MemoryKV {
34    manager: Arc<MemoryTxnManager>,
35}
36
37impl MemoryKV {
38    /// Creates a new, purely transient in-memory KV store.
39    pub fn new() -> Self {
40        Self {
41            manager: Arc::new(MemoryTxnManager::new(None, None, None)),
42        }
43    }
44
45    /// Returns current in-memory usage statistics.
46    pub fn memory_stats(&self) -> MemoryStats {
47        self.manager.memory_stats()
48    }
49
50    /// Creates a new in-memory KV store with an optional memory limit.
51    pub fn new_with_limit(limit: Option<usize>) -> Self {
52        Self {
53            manager: Arc::new(MemoryTxnManager::new_with_limit(limit)),
54        }
55    }
56
57    /// Opens a persistent in-memory KV store from a file path.
58    pub fn open(path: &Path) -> Result<Self> {
59        let wal_writer = WalWriter::new(path)?;
60        let sstable_path = path.with_extension("sst");
61        let manager = Arc::new(MemoryTxnManager::new(
62            Some(wal_writer),
63            Some(path.to_path_buf()),
64            Some(sstable_path),
65        ));
66        manager.recover()?;
67        Ok(Self { manager })
68    }
69
70    /// Opens a persistent KV store with I/O hooks (test only).
71    #[cfg(feature = "test-hooks")]
72    pub fn open_with_io_hooks(path: &Path, hooks: Arc<dyn IoHooks>) -> Result<Self> {
73        let wal_writer = WalWriter::new(path)?;
74        let sstable_path = path.with_extension("sst");
75        let manager = Arc::new(MemoryTxnManager::new(
76            Some(wal_writer),
77            Some(path.to_path_buf()),
78            Some(sstable_path),
79        ));
80        manager.set_io_hooks(Some(hooks));
81        manager.recover()?;
82        Ok(Self { manager })
83    }
84
85    /// Opens a persistent KV store with crash hooks (test only).
86    #[cfg(feature = "test-hooks")]
87    pub fn open_with_crash_hooks(path: &Path, crash_sim: Arc<CrashSimulator>) -> Result<Self> {
88        let wal_writer = WalWriter::new(path)?;
89        let sstable_path = path.with_extension("sst");
90        let manager = Arc::new(MemoryTxnManager::new(
91            Some(wal_writer),
92            Some(path.to_path_buf()),
93            Some(sstable_path),
94        ));
95        manager.set_crash_sim(Some(crash_sim));
96        manager.recover()?;
97        Ok(Self { manager })
98    }
99
100    /// Flushes the in-memory data to an SSTable.
101    pub fn flush(&self) -> Result<()> {
102        self.manager.flush()
103    }
104}
105
106impl Default for MemoryKV {
107    fn default() -> Self {
108        Self::new()
109    }
110}
111
112impl KVStore for MemoryKV {
113    type Transaction<'a> = MemoryTransaction<'a>;
114    type Manager<'a> = &'a MemoryTxnManager;
115
116    fn txn_manager(&self) -> Self::Manager<'_> {
117        &self.manager
118    }
119
120    fn begin(&self, mode: TxnMode) -> Result<Self::Transaction<'_>> {
121        self.manager.begin_internal(mode)
122    }
123
124    fn runtime_stats(&self) -> Option<crate::kv::RuntimeStats> {
125        Some(crate::kv::RuntimeStats::Memory(self.memory_stats()))
126    }
127
128    fn set_memory_limit_bytes(&self, limit: Option<usize>) -> Result<()> {
129        self.manager.set_memory_limit(limit);
130        Ok(())
131    }
132}
133
134impl OwnedKVStore for MemoryKV {
135    fn begin_owned_kv_transaction(
136        self: Arc<Self>,
137        mode: TxnMode,
138    ) -> Result<Box<dyn OwnedKVTransaction>> {
139        Ok(Box::new(OwnedMemoryTransaction::new(
140            self.manager.clone(),
141            mode,
142        )))
143    }
144}
145
146// The internal value stored in the BTreeMap, containing the data and its version.
147type VersionedValue = (Value, u64);
148
149/// The underlying shared state for the in-memory store.
150struct MemorySharedState {
151    /// The main data store, mapping keys to versioned values.
152    data: RwLock<BTreeMap<Key, VersionedValue>>,
153    /// The next transaction ID to be allocated.
154    next_txn_id: AtomicU64,
155    /// The current commit version of the database. Incremented on every successful commit.
156    commit_version: AtomicU64,
157    /// The WAL writer. If None, the store is transient.
158    wal_writer: Option<RwLock<WalWriter>>,
159    /// Optional WAL path for replay on reopen.
160    wal_path: Option<PathBuf>,
161    /// Optional SSTable reader for read-through.
162    sstable: RwLock<Option<SstableReader>>,
163    /// Optional SSTable path for flush/reopen.
164    sstable_path: Option<PathBuf>,
165    /// Optional memory upper limit (bytes) for in-memory mode。
166    memory_limit: RwLock<Option<usize>>,
167    /// Current memory consumption (bytes) tracked across operations。
168    current_memory: AtomicUsize,
169    /// Safe snapshot gate for owned incremental cursors.  It prevents writers from replacing
170    /// the map while a cursor owns its stable read view, without self-referential lock guards.
171    owned_snapshot_gate: Arc<OwnedSnapshotGate>,
172    #[cfg(feature = "test-hooks")]
173    /// Optional I/O hooks for fault injection (test only).
174    io_hooks: RwLock<Option<Arc<dyn IoHooks>>>,
175    #[cfg(feature = "test-hooks")]
176    /// Optional crash simulator for test-only crash points.
177    crash_sim: RwLock<Option<Arc<CrashSimulator>>>,
178}
179
180#[derive(Default)]
181struct OwnedSnapshotGate {
182    state: Mutex<OwnedSnapshotGateState>,
183    changed: Condvar,
184}
185
186#[derive(Default)]
187struct OwnedSnapshotGateState {
188    readers: usize,
189    writer: bool,
190}
191
192impl OwnedSnapshotGate {
193    fn acquire_reader(self: &Arc<Self>) -> OwnedSnapshotReader {
194        let mut state = self
195            .state
196            .lock()
197            .expect("owned snapshot gate mutex poisoned");
198        while state.writer {
199            state = self
200                .changed
201                .wait(state)
202                .expect("owned snapshot gate mutex poisoned");
203        }
204        state.readers = state.readers.saturating_add(1);
205        OwnedSnapshotReader {
206            gate: self.clone(),
207            released: false,
208        }
209    }
210
211    fn acquire_writer(self: &Arc<Self>) -> OwnedSnapshotWriter {
212        let mut state = self
213            .state
214            .lock()
215            .expect("owned snapshot gate mutex poisoned");
216        while state.writer || state.readers != 0 {
217            state = self
218                .changed
219                .wait(state)
220                .expect("owned snapshot gate mutex poisoned");
221        }
222        state.writer = true;
223        OwnedSnapshotWriter {
224            gate: self.clone(),
225            released: false,
226        }
227    }
228}
229
230struct OwnedSnapshotReader {
231    gate: Arc<OwnedSnapshotGate>,
232    released: bool,
233}
234
235impl Drop for OwnedSnapshotReader {
236    fn drop(&mut self) {
237        if !self.released {
238            let mut state = self
239                .gate
240                .state
241                .lock()
242                .expect("owned snapshot gate mutex poisoned");
243            state.readers = state.readers.saturating_sub(1);
244            self.released = true;
245            self.gate.changed.notify_all();
246        }
247    }
248}
249
250struct OwnedSnapshotWriter {
251    gate: Arc<OwnedSnapshotGate>,
252    released: bool,
253}
254
255impl Drop for OwnedSnapshotWriter {
256    fn drop(&mut self) {
257        if !self.released {
258            let mut state = self
259                .gate
260                .state
261                .lock()
262                .expect("owned snapshot gate mutex poisoned");
263            state.writer = false;
264            self.released = true;
265            self.gate.changed.notify_all();
266        }
267    }
268}
269
270impl MemorySharedState {
271    /// Check whether adding `additional` bytes would exceed the memory limit.
272    fn check_memory_limit(&self, additional: usize) -> Result<()> {
273        if let Some(limit) = *self.memory_limit.read().unwrap() {
274            let current = self.current_memory.load(Ordering::Relaxed);
275            let requested = current.saturating_add(additional);
276            if requested > limit {
277                return Err(Error::MemoryLimitExceeded { limit, requested });
278            }
279        }
280        Ok(())
281    }
282
283    /// Return current memory usage statistics.
284    fn memory_stats(&self) -> MemoryStats {
285        let kv_bytes = self.current_memory.load(Ordering::Relaxed);
286        MemoryStats {
287            total_bytes: kv_bytes,
288            kv_bytes,
289            index_bytes: 0,
290        }
291    }
292
293    /// Recompute tracked memory usage from existing data (used after recovery).
294    fn recompute_current_memory(&self) {
295        let data = self.data.read().unwrap();
296        let mut total = 0usize;
297        for (k, (v, _)) in data.iter() {
298            total = total.saturating_add(k.len() + v.len());
299        }
300        self.current_memory.store(total, Ordering::Relaxed);
301    }
302}
303
304/// A transaction manager backed by an in-memory map and optional WAL.
305pub struct MemoryTxnManager {
306    state: Arc<MemorySharedState>,
307}
308
309impl MemoryTxnManager {
310    fn new_with_params(
311        wal_writer: Option<WalWriter>,
312        wal_path: Option<PathBuf>,
313        sstable_path: Option<PathBuf>,
314        memory_limit: Option<usize>,
315    ) -> Self {
316        Self {
317            state: Arc::new(MemorySharedState {
318                data: RwLock::new(BTreeMap::new()),
319                next_txn_id: AtomicU64::new(1),
320                commit_version: AtomicU64::new(0),
321                wal_writer: wal_writer.map(RwLock::new),
322                wal_path,
323                sstable: RwLock::new(None),
324                sstable_path,
325                memory_limit: RwLock::new(memory_limit),
326                current_memory: AtomicUsize::new(0),
327                owned_snapshot_gate: Arc::new(OwnedSnapshotGate::default()),
328                #[cfg(feature = "test-hooks")]
329                io_hooks: RwLock::new(None),
330                #[cfg(feature = "test-hooks")]
331                crash_sim: RwLock::new(None),
332            }),
333        }
334    }
335
336    fn new(
337        wal_writer: Option<WalWriter>,
338        wal_path: Option<PathBuf>,
339        sstable_path: Option<PathBuf>,
340    ) -> Self {
341        Self::new_with_params(wal_writer, wal_path, sstable_path, None)
342    }
343
344    /// Creates an in-memory manager with an optional memory limit.
345    pub fn new_with_limit(limit: Option<usize>) -> Self {
346        Self::new_with_params(None, None, None, limit)
347    }
348
349    #[cfg(feature = "test-hooks")]
350    fn set_io_hooks(&self, hooks: Option<Arc<dyn IoHooks>>) {
351        let mut guard = self.state.io_hooks.write().unwrap();
352        *guard = hooks;
353    }
354
355    #[cfg(feature = "test-hooks")]
356    fn set_crash_sim(&self, crash_sim: Option<Arc<CrashSimulator>>) {
357        let mut guard = self.state.crash_sim.write().unwrap();
358        *guard = crash_sim;
359    }
360
361    #[cfg(feature = "test-hooks")]
362    fn io_hooks(&self) -> Option<Arc<dyn IoHooks>> {
363        self.state.io_hooks.read().unwrap().clone()
364    }
365
366    #[cfg(feature = "test-hooks")]
367    fn crash_sim(&self) -> Option<Arc<CrashSimulator>> {
368        self.state.crash_sim.read().unwrap().clone()
369    }
370
371    /// Returns current memory usage statistics.
372    pub fn memory_stats(&self) -> MemoryStats {
373        self.state.memory_stats()
374    }
375
376    /// Update the configured memory limit at runtime.
377    pub fn set_memory_limit(&self, limit: Option<usize>) {
378        let mut guard = self.state.memory_limit.write().unwrap();
379        *guard = limit;
380    }
381
382    /// Returns a snapshot clone of all key/value pairs.
383    pub fn snapshot(&self) -> Vec<(Key, Value)> {
384        let data = self.state.data.read().unwrap();
385        data.iter()
386            .map(|(k, (v, _))| (k.clone(), v.clone()))
387            .collect()
388    }
389
390    /// Clears all data and resets memory accounting.
391    pub fn clear_all(&self) {
392        let _snapshot_writer = self.state.owned_snapshot_gate.acquire_writer();
393        let mut data = self.state.data.write().unwrap();
394        data.clear();
395        drop(data);
396        self.state.current_memory.store(0, Ordering::Relaxed);
397        self.state.commit_version.store(0, Ordering::Relaxed);
398    }
399
400    /// Runs compaction if it can fit within the configured memory limit.
401    /// Returns Ok(true) when compaction executed, Ok(false) when skipped.
402    pub fn compact_with_limit<F>(
403        &self,
404        input_bytes: usize,
405        output_bytes: usize,
406        run: F,
407    ) -> Result<bool>
408    where
409        F: FnOnce() -> Result<()>,
410    {
411        if let Some(limit) = *self.state.memory_limit.read().unwrap() {
412            let current = self.state.current_memory.load(Ordering::Relaxed);
413            // predicted usage after compaction: current - input + output (clamped at 0)
414            let prospective = current
415                .saturating_sub(input_bytes)
416                .saturating_add(output_bytes);
417            if prospective > limit {
418                warn!(
419                    limit,
420                    requested = prospective,
421                    "compaction skipped due to memory limit"
422                );
423                return Ok(false);
424            }
425        }
426
427        run()?;
428
429        // Update tracked memory to reflect compaction result.
430        let current = self.state.current_memory.load(Ordering::Relaxed);
431        let new_usage = current
432            .saturating_sub(input_bytes)
433            .saturating_add(output_bytes);
434        self.state
435            .current_memory
436            .store(new_usage, Ordering::Relaxed);
437        Ok(true)
438    }
439
440    #[cfg(feature = "test-hooks")]
441    fn trigger_crash(&self, operation: CrashOperation, timing: CrashTiming) {
442        if let Some(sim) = self.crash_sim() {
443            sim.check_crash(operation, timing);
444        }
445    }
446
447    #[cfg(feature = "test-hooks")]
448    fn notify_wal_hooks(&self, data: &[u8], timing: CrashTiming) -> Result<()> {
449        match timing {
450            CrashTiming::Before => {
451                self.trigger_crash(CrashOperation::WalWrite, CrashTiming::Before);
452                if let Some(hooks) = self.io_hooks() {
453                    hooks.before_wal_write(data).map_err(Error::Io)?;
454                    hooks.before_fsync().map_err(Error::Io)?;
455                }
456                self.trigger_crash(CrashOperation::WalFsync, CrashTiming::Before);
457            }
458            CrashTiming::During => {
459                self.trigger_crash(CrashOperation::WalWrite, CrashTiming::During);
460                self.trigger_crash(CrashOperation::WalFsync, CrashTiming::During);
461            }
462            CrashTiming::After => {
463                self.trigger_crash(CrashOperation::WalWrite, CrashTiming::After);
464                self.trigger_crash(CrashOperation::WalFsync, CrashTiming::After);
465                if let Some(hooks) = self.io_hooks() {
466                    hooks.after_wal_write(data).map_err(Error::Io)?;
467                    hooks.after_fsync().map_err(Error::Io)?;
468                }
469            }
470        }
471        Ok(())
472    }
473
474    #[cfg(feature = "test-hooks")]
475    fn notify_compaction(&self, timing: CrashTiming) {
476        self.trigger_crash(CrashOperation::Compaction, timing);
477        if let Some(hooks) = self.io_hooks() {
478            match timing {
479                CrashTiming::Before => hooks.on_compaction_start(),
480                CrashTiming::After => hooks.on_compaction_end(),
481                CrashTiming::During => {}
482            }
483        }
484    }
485
486    fn append_wal_record(&self, wal: &mut WalWriter, record: &WalRecord) -> Result<()> {
487        #[cfg(feature = "test-hooks")]
488        {
489            let data =
490                bincode::serialize(record).map_err(|e| Error::Io(std::io::Error::other(e)))?;
491            self.notify_wal_hooks(&data, CrashTiming::Before)?;
492            self.notify_wal_hooks(&data, CrashTiming::During)?;
493            wal.append(record)?;
494            self.notify_wal_hooks(&data, CrashTiming::After)?;
495            Ok(())
496        }
497
498        #[cfg(not(feature = "test-hooks"))]
499        {
500            wal.append(record)
501        }
502    }
503
504    fn write_wal(&self, txn_id: TxnId, writes: &BTreeMap<Key, Option<Value>>) -> Result<()> {
505        if let Some(wal_lock) = &self.state.wal_writer {
506            let mut wal = wal_lock.write().unwrap();
507            self.append_wal_record(&mut wal, &WalRecord::Begin(txn_id))?;
508            for (key, value) in writes {
509                let record = match value {
510                    Some(v) => WalRecord::Put(txn_id, key.clone(), v.clone()),
511                    None => WalRecord::Delete(txn_id, key.clone()),
512                };
513                self.append_wal_record(&mut wal, &record)?;
514            }
515            self.append_wal_record(&mut wal, &WalRecord::Commit(txn_id))?;
516            // Durability point: fsync once at the commit boundary so the whole
517            // transaction (Begin..Commit) reaches stable storage together. The
518            // commit is only acknowledged after this returns (CORE-5.1).
519            wal.sync()?;
520        }
521        Ok(())
522    }
523
524    /// In-memory compaction entrypoint that rebuilds the map while honoring memory limits.
525    pub fn compact_in_memory(&self) -> Result<bool> {
526        #[cfg(feature = "test-hooks")]
527        self.notify_compaction(CrashTiming::Before);
528
529        let snapshot_bytes = {
530            let data = self.state.data.read().unwrap();
531            let mut bytes = 0usize;
532            for (k, (v, _)) in data.iter() {
533                bytes = bytes.saturating_add(k.len() + v.len());
534            }
535            bytes
536        };
537
538        let executed = self.compact_with_limit(snapshot_bytes, snapshot_bytes, || {
539            let data = self.state.data.read().unwrap();
540            let mut rebuilt = BTreeMap::new();
541            for (k, (v, version)) in data.iter() {
542                rebuilt.insert(k.clone(), (v.clone(), *version));
543            }
544            drop(data);
545
546            #[cfg(feature = "test-hooks")]
547            self.notify_compaction(CrashTiming::During);
548
549            let _snapshot_writer = self.state.owned_snapshot_gate.acquire_writer();
550            let mut write_guard = self.state.data.write().unwrap();
551            *write_guard = rebuilt;
552            Ok(())
553        })?;
554
555        #[cfg(feature = "test-hooks")]
556        self.notify_compaction(CrashTiming::After);
557
558        Ok(executed)
559    }
560
561    /// Flushes the current in-memory data to an SSTable file.
562    pub fn flush(&self) -> Result<()> {
563        let Some(path) = self.state.sstable_path.as_ref() else {
564            return Ok(());
565        };
566
567        #[cfg(feature = "test-hooks")]
568        self.notify_compaction(CrashTiming::Before);
569
570        let data = self.state.data.read().unwrap();
571        let mut writer = SstableWriter::create(path)?;
572        for (key, (value, _version)) in data.iter() {
573            #[cfg(feature = "test-hooks")]
574            {
575                let mut record = Vec::with_capacity(key.len() + value.len());
576                record.extend_from_slice(key);
577                record.extend_from_slice(value);
578                self.trigger_crash(CrashOperation::SstWrite, CrashTiming::Before);
579                if let Some(hooks) = self.io_hooks() {
580                    hooks.before_sst_write(&record).map_err(Error::Io)?;
581                }
582                self.trigger_crash(CrashOperation::SstWrite, CrashTiming::During);
583            }
584
585            writer.append(key, value)?;
586
587            #[cfg(feature = "test-hooks")]
588            self.trigger_crash(CrashOperation::SstWrite, CrashTiming::After);
589        }
590        drop(data);
591
592        #[cfg(feature = "test-hooks")]
593        self.trigger_crash(CrashOperation::SstFinalize, CrashTiming::Before);
594        let _footer = writer.finish()?;
595        #[cfg(feature = "test-hooks")]
596        self.trigger_crash(CrashOperation::SstFinalize, CrashTiming::After);
597        let reader = SstableReader::open(path)?;
598        // Also emit a placeholder vector segment alongside SSTable for future vector recovery.
599        let vec_path = path.with_extension("vec");
600        write_empty_vector_segment(&vec_path)?;
601
602        let mut slot = self.state.sstable.write().unwrap();
603        *slot = Some(reader);
604
605        #[cfg(feature = "test-hooks")]
606        self.notify_compaction(CrashTiming::After);
607        Ok(())
608    }
609
610    /// Replays the WAL to restore the state of the in-memory map.
611    fn replay(&self) -> Result<()> {
612        let path = match &self.state.wal_path {
613            Some(p) => p,
614            None => return Ok(()),
615        };
616        if !path.exists() || std::fs::metadata(path)?.len() == 0 {
617            return Ok(());
618        }
619
620        let _snapshot_writer = self.state.owned_snapshot_gate.acquire_writer();
621        let mut data = self.state.data.write().unwrap();
622        let mut max_txn_id = 0;
623        let mut max_version = self.state.commit_version.load(Ordering::Acquire);
624        let reader = WalReader::new(path)?;
625        let mut pending_txns: HashMap<TxnId, Vec<(Key, Option<Value>)>> = HashMap::new();
626
627        for record_result in reader {
628            match record_result? {
629                WalRecord::Begin(txn_id) => {
630                    max_txn_id = max_txn_id.max(txn_id.0);
631                    pending_txns.entry(txn_id).or_default();
632                }
633                WalRecord::Put(txn_id, key, value) => {
634                    max_txn_id = max_txn_id.max(txn_id.0);
635                    pending_txns
636                        .entry(txn_id)
637                        .or_default()
638                        .push((key, Some(value)));
639                }
640                WalRecord::Delete(txn_id, key) => {
641                    max_txn_id = max_txn_id.max(txn_id.0);
642                    pending_txns.entry(txn_id).or_default().push((key, None));
643                }
644                WalRecord::Commit(txn_id) => {
645                    if let Some(writes) = pending_txns.remove(&txn_id) {
646                        max_version += 1;
647                        for (key, value) in writes {
648                            if let Some(v) = value {
649                                data.insert(key, (v, max_version));
650                            } else {
651                                data.remove(&key);
652                            }
653                        }
654                    }
655                }
656            }
657        }
658
659        self.state
660            .next_txn_id
661            .store(max_txn_id + 1, Ordering::SeqCst);
662        self.state
663            .commit_version
664            .store(max_version, Ordering::SeqCst);
665        Ok(())
666    }
667
668    fn load_sstable(&self) -> Result<()> {
669        let path = match &self.state.sstable_path {
670            Some(p) => p,
671            None => return Ok(()),
672        };
673        if !path.exists() {
674            return Ok(());
675        }
676
677        let mut reader = match SstableReader::open(path) {
678            Ok(reader) => reader,
679            // Crash-recovery tolerance: a truncated or corrupt SSTable means the
680            // durable segment was torn by a crash and its contents are
681            // unrecoverable. Discard the unreadable segment and continue recovery
682            // from the WAL — this mirrors `WalReader`'s torn-tail handling. Only
683            // corruption signatures are tolerated; genuine I/O faults still
684            // propagate so real failures are not masked.
685            Err(e @ (Error::InvalidFormat(_) | Error::ChecksumMismatch)) => {
686                warn!(
687                    path = %path.display(),
688                    error = %e,
689                    "discarding unreadable SSTable during recovery; replaying WAL only"
690                );
691                return Ok(());
692            }
693            Err(Error::Io(io)) if io.kind() == std::io::ErrorKind::UnexpectedEof => {
694                warn!(
695                    path = %path.display(),
696                    "discarding truncated SSTable during recovery; replaying WAL only"
697                );
698                return Ok(());
699            }
700            Err(e) => return Err(e),
701        };
702        let mut data = self.state.data.write().unwrap();
703        let mut version = self.state.commit_version.load(Ordering::Acquire);
704
705        let keys: Vec<Key> = reader
706            .index()
707            .iter()
708            .map(|entry| entry.key.clone())
709            .collect();
710
711        for key in keys {
712            if let Some(value) = reader.get(&key)? {
713                version += 1;
714                data.insert(key, (value, version));
715            }
716        }
717
718        self.state.commit_version.store(version, Ordering::SeqCst);
719        let mut slot = self.state.sstable.write().unwrap();
720        *slot = Some(reader);
721        Ok(())
722    }
723
724    /// Loads SSTable then replays WAL to restore state.
725    fn recover(&self) -> Result<()> {
726        self.load_sstable()?;
727        self.replay()?;
728        self.state.recompute_current_memory();
729        Ok(())
730    }
731
732    fn sstable_get(&self, key: &Key) -> Result<Option<Value>> {
733        let mut guard = self.state.sstable.write().unwrap();
734        if let Some(reader) = guard.as_mut() {
735            return reader.get(key);
736        }
737        Ok(None)
738    }
739
740    fn begin_internal(&self, mode: TxnMode) -> Result<MemoryTransaction<'_>> {
741        let txn_id = self.state.next_txn_id.fetch_add(1, Ordering::SeqCst);
742        let start_version = self.state.commit_version.load(Ordering::Acquire);
743        Ok(MemoryTransaction::new(
744            self,
745            TxnId(txn_id),
746            mode,
747            start_version,
748        ))
749    }
750}
751
752impl<'a> TxnManager<'a, MemoryTransaction<'a>> for &'a MemoryTxnManager {
753    fn begin(&'a self, mode: TxnMode) -> Result<MemoryTransaction<'a>> {
754        self.begin_internal(mode)
755    }
756
757    fn commit(&'a self, mut txn: MemoryTransaction<'a>) -> Result<()> {
758        if txn.state != TxnState::Active {
759            return Err(Error::TxnClosed);
760        }
761        if txn.mode == TxnMode::ReadOnly || txn.writes.is_empty() {
762            txn.state = TxnState::Committed;
763            return Ok(());
764        }
765
766        let _snapshot_writer = self.state.owned_snapshot_gate.acquire_writer();
767        let mut data = self.state.data.write().unwrap();
768
769        for key in txn.read_set.keys() {
770            let current_version = data.get(key).map(|(_, v)| *v).unwrap_or(0);
771            if current_version > txn.start_version {
772                return Err(Error::TxnConflict);
773            }
774        }
775
776        // Detect write-write conflicts even when the key was never read.
777        for key in txn.writes.keys() {
778            let current_version = data.get(key).map(|(_, v)| *v).unwrap_or(0);
779            if current_version > txn.start_version {
780                return Err(Error::TxnConflict);
781            }
782        }
783
784        // Compute prospective memory usage and enforce limits before mutating state.
785        let mut delta: isize = 0;
786        for (key, value) in &txn.writes {
787            let current_size = data.get(key).map(|(v, _)| key.len() + v.len()).unwrap_or(0);
788            let new_size = match value {
789                Some(v) => key.len() + v.len(),
790                None => 0,
791            };
792            delta += new_size as isize - current_size as isize;
793        }
794
795        let current_mem = self.state.current_memory.load(Ordering::Relaxed);
796        let prospective = if delta >= 0 {
797            current_mem.saturating_add(delta as usize)
798        } else {
799            current_mem.saturating_sub(delta.unsigned_abs())
800        };
801
802        if delta > 0 {
803            self.state.check_memory_limit(delta as usize)?;
804        }
805
806        let commit_version = self.state.commit_version.fetch_add(1, Ordering::AcqRel) + 1;
807
808        self.write_wal(txn.id, &txn.writes)?;
809
810        for (key, value) in std::mem::take(&mut txn.writes) {
811            if let Some(v) = value {
812                data.insert(key, (v, commit_version));
813            } else {
814                data.remove(&key);
815            }
816        }
817
818        self.state
819            .current_memory
820            .store(prospective, Ordering::Relaxed);
821
822        txn.state = TxnState::Committed;
823        Ok(())
824    }
825
826    fn rollback(&'a self, mut txn: MemoryTransaction<'a>) -> Result<()> {
827        if txn.state != TxnState::Active {
828            return Err(Error::TxnClosed);
829        }
830        txn.state = TxnState::RolledBack;
831        Ok(())
832    }
833}
834
835/// An in-memory transaction that enforces snapshot isolation.
836pub struct MemoryTransaction<'a> {
837    manager: &'a MemoryTxnManager,
838    id: TxnId,
839    mode: TxnMode,
840    state: TxnState,
841    start_version: u64,
842    writes: BTreeMap<Key, Option<Value>>,
843    read_set: HashMap<Key, u64>,
844}
845
846impl<'a> MemoryTransaction<'a> {
847    fn new(manager: &'a MemoryTxnManager, id: TxnId, mode: TxnMode, start_version: u64) -> Self {
848        Self {
849            manager,
850            id,
851            mode,
852            state: TxnState::Active,
853            start_version,
854            writes: BTreeMap::new(),
855            read_set: HashMap::new(),
856        }
857    }
858
859    fn ensure_active(&self) -> Result<()> {
860        if self.state != TxnState::Active {
861            return Err(Error::TxnClosed);
862        }
863        Ok(())
864    }
865
866    /// トランザクションを消費せずにロールバックする。
867    pub(crate) fn rollback_in_place(&mut self) -> Result<()> {
868        if self.state != TxnState::Active {
869            return Err(Error::TxnClosed);
870        }
871        self.state = TxnState::RolledBack;
872        Ok(())
873    }
874
875    fn scan_range_internal(&mut self, start: &[u8], end: &[u8]) -> MergedScanIter<'_> {
876        let start_vec = start.to_vec();
877        let end_vec = end.to_vec();
878        let data_guard = self.manager.state.data.read().unwrap();
879        let data_ptr: *const BTreeMap<Key, VersionedValue> = &*data_guard;
880        let data_iter = unsafe {
881            // Safety: data_guard keeps the map alive for the lifetime of the iterator.
882            (&*data_ptr).range((Included(start_vec.clone()), Excluded(end_vec.clone())))
883        };
884        let write_iter = self
885            .writes
886            .range((Included(start_vec.clone()), Excluded(end_vec.clone())));
887
888        MergedScanIter::new(
889            data_guard,
890            data_iter,
891            write_iter,
892            None,
893            Some(end_vec),
894            self.start_version,
895            &mut self.read_set,
896        )
897    }
898
899    fn scan_prefix_internal(&mut self, prefix: &[u8]) -> MergedScanIter<'_> {
900        let prefix_vec = prefix.to_vec();
901        let data_guard = self.manager.state.data.read().unwrap();
902        let data_ptr: *const BTreeMap<Key, VersionedValue> = &*data_guard;
903        let data_iter = unsafe {
904            // Safety: data_guard keeps the map alive for the lifetime of the iterator.
905            (&*data_ptr).range(prefix_vec.clone()..)
906        };
907        let write_iter = self.writes.range(prefix_vec.clone()..);
908        MergedScanIter::new(
909            data_guard,
910            data_iter,
911            write_iter,
912            Some(prefix_vec),
913            None,
914            self.start_version,
915            &mut self.read_set,
916        )
917    }
918}
919
920impl<'a> KVTransaction<'a> for MemoryTransaction<'a> {
921    fn id(&self) -> TxnId {
922        self.id
923    }
924
925    fn mode(&self) -> TxnMode {
926        self.mode
927    }
928
929    fn get(&mut self, key: &Key) -> Result<Option<Value>> {
930        if self.state != TxnState::Active {
931            return Err(Error::TxnClosed);
932        }
933
934        if let Some(value) = self.writes.get(key) {
935            return Ok(value.clone());
936        }
937
938        let result = {
939            let data = self.manager.state.data.read().unwrap();
940            data.get(key).cloned()
941        };
942
943        if let Some((v, version)) = result {
944            self.read_set.insert(key.clone(), version);
945            return Ok(Some(v));
946        }
947
948        // Read-through to SSTable if not found in memory.
949        if let Some(value) = self.manager.sstable_get(key)? {
950            let version = self.manager.state.commit_version.load(Ordering::Acquire);
951            self.read_set.insert(key.clone(), version);
952            return Ok(Some(value));
953        }
954
955        Ok(None)
956    }
957
958    fn put(&mut self, key: Key, value: Value) -> Result<()> {
959        if self.state != TxnState::Active {
960            return Err(Error::TxnClosed);
961        }
962        if self.mode == TxnMode::ReadOnly {
963            return Err(Error::TxnReadOnly);
964        }
965        self.writes.insert(key, Some(value));
966        Ok(())
967    }
968
969    fn delete(&mut self, key: Key) -> Result<()> {
970        if self.state != TxnState::Active {
971            return Err(Error::TxnClosed);
972        }
973        if self.mode == TxnMode::ReadOnly {
974            return Err(Error::TxnReadOnly);
975        }
976        self.writes.insert(key, None);
977        Ok(())
978    }
979
980    fn scan_prefix(
981        &mut self,
982        prefix: &[u8],
983    ) -> Result<Box<dyn Iterator<Item = (Key, Value)> + '_>> {
984        self.ensure_active()?;
985        let iter = self
986            .scan_prefix_internal(prefix)
987            .filter_map(|(k, v)| v.map(|val| (k, val)));
988        Ok(Box::new(iter))
989    }
990
991    fn scan_range(
992        &mut self,
993        start: &[u8],
994        end: &[u8],
995    ) -> Result<Box<dyn Iterator<Item = (Key, Value)> + '_>> {
996        self.ensure_active()?;
997        let iter = self
998            .scan_range_internal(start, end)
999            .filter_map(|(k, v)| v.map(|val| (k, val)));
1000        Ok(Box::new(iter))
1001    }
1002
1003    fn commit_self(mut self) -> Result<()> {
1004        if self.state != TxnState::Active {
1005            return Err(Error::TxnClosed);
1006        }
1007        if self.mode == TxnMode::ReadOnly || self.writes.is_empty() {
1008            self.state = TxnState::Committed;
1009            return Ok(());
1010        }
1011
1012        let _snapshot_writer = self.manager.state.owned_snapshot_gate.acquire_writer();
1013        let mut data = self.manager.state.data.write().unwrap();
1014
1015        // Check read-set for conflicts
1016        for key in self.read_set.keys() {
1017            let current_version = data.get(key).map(|(_, v)| *v).unwrap_or(0);
1018            if current_version > self.start_version {
1019                return Err(Error::TxnConflict);
1020            }
1021        }
1022
1023        // Check write-write conflicts
1024        for key in self.writes.keys() {
1025            let current_version = data.get(key).map(|(_, v)| *v).unwrap_or(0);
1026            if current_version > self.start_version {
1027                return Err(Error::TxnConflict);
1028            }
1029        }
1030
1031        // Compute prospective memory usage
1032        let mut delta: isize = 0;
1033        for (key, value) in &self.writes {
1034            let current_size = data.get(key).map(|(v, _)| key.len() + v.len()).unwrap_or(0);
1035            let new_size = match value {
1036                Some(v) => key.len() + v.len(),
1037                None => 0,
1038            };
1039            delta += new_size as isize - current_size as isize;
1040        }
1041
1042        let current_mem = self.manager.state.current_memory.load(Ordering::Relaxed);
1043        let prospective = if delta >= 0 {
1044            current_mem.saturating_add(delta as usize)
1045        } else {
1046            current_mem.saturating_sub(delta.unsigned_abs())
1047        };
1048
1049        if delta > 0 {
1050            self.manager.state.check_memory_limit(delta as usize)?;
1051        }
1052
1053        let commit_version = self
1054            .manager
1055            .state
1056            .commit_version
1057            .fetch_add(1, Ordering::AcqRel)
1058            + 1;
1059
1060        // WAL write
1061        self.manager.write_wal(self.id, &self.writes)?;
1062
1063        // Apply writes
1064        for (key, value) in std::mem::take(&mut self.writes) {
1065            if let Some(v) = value {
1066                data.insert(key, (v, commit_version));
1067            } else {
1068                data.remove(&key);
1069            }
1070        }
1071
1072        self.manager
1073            .state
1074            .current_memory
1075            .store(prospective, Ordering::Relaxed);
1076
1077        self.state = TxnState::Committed;
1078        Ok(())
1079    }
1080
1081    fn rollback_self(mut self) -> Result<()> {
1082        if self.state != TxnState::Active {
1083            return Err(Error::TxnClosed);
1084        }
1085        self.state = TxnState::RolledBack;
1086        Ok(())
1087    }
1088}
1089
1090/// An owned MemoryKV transaction whose manager and mutable snapshot state both outlive callers.
1091struct OwnedMemoryTransaction {
1092    manager: Arc<MemoryTxnManager>,
1093    state: Arc<Mutex<OwnedMemoryTransactionState>>,
1094}
1095
1096struct OwnedMemoryTransactionState {
1097    id: TxnId,
1098    mode: TxnMode,
1099    state: TxnState,
1100    start_version: u64,
1101    writes: BTreeMap<Key, Option<Value>>,
1102    read_set: HashMap<Key, u64>,
1103    cursor_open: bool,
1104}
1105
1106impl OwnedMemoryTransaction {
1107    fn new(manager: Arc<MemoryTxnManager>, mode: TxnMode) -> Self {
1108        let id = TxnId(manager.state.next_txn_id.fetch_add(1, Ordering::SeqCst));
1109        let start_version = manager.state.commit_version.load(Ordering::Acquire);
1110        Self {
1111            manager,
1112            state: Arc::new(Mutex::new(OwnedMemoryTransactionState {
1113                id,
1114                mode,
1115                state: TxnState::Active,
1116                start_version,
1117                writes: BTreeMap::new(),
1118                read_set: HashMap::new(),
1119                cursor_open: false,
1120            })),
1121        }
1122    }
1123
1124    fn open_cursor(
1125        &mut self,
1126        start: Option<Key>,
1127        prefix: Option<Vec<u8>>,
1128        end: Option<Key>,
1129    ) -> Result<Box<dyn OwnedKVScan>> {
1130        let snapshot = self.manager.state.owned_snapshot_gate.acquire_reader();
1131        let mut state = self
1132            .state
1133            .lock()
1134            .expect("owned memory transaction mutex poisoned");
1135        if state.state != TxnState::Active || state.cursor_open {
1136            return Err(Error::TxnClosed);
1137        }
1138        state.cursor_open = true;
1139        drop(state);
1140        Ok(Box::new(OwnedMemoryCursor {
1141            manager: self.manager.clone(),
1142            transaction: self.state.clone(),
1143            snapshot: Some(snapshot),
1144            last_key: None,
1145            start,
1146            prefix,
1147            end,
1148        }))
1149    }
1150
1151    fn ensure_active(state: &OwnedMemoryTransactionState) -> Result<()> {
1152        if state.state != TxnState::Active {
1153            return Err(Error::TxnClosed);
1154        }
1155        Ok(())
1156    }
1157}
1158
1159impl OwnedKVTransaction for OwnedMemoryTransaction {
1160    fn id(&self) -> TxnId {
1161        self.state
1162            .lock()
1163            .expect("owned memory transaction mutex poisoned")
1164            .id
1165    }
1166
1167    fn mode(&self) -> TxnMode {
1168        self.state
1169            .lock()
1170            .expect("owned memory transaction mutex poisoned")
1171            .mode
1172    }
1173
1174    fn get(&mut self, key: &Key) -> Result<Option<Value>> {
1175        let mut state = self
1176            .state
1177            .lock()
1178            .expect("owned memory transaction mutex poisoned");
1179        Self::ensure_active(&state)?;
1180        if let Some(value) = state.writes.get(key) {
1181            return Ok(value.clone());
1182        }
1183
1184        let result = {
1185            let data = self.manager.state.data.read().unwrap();
1186            data.get(key).cloned()
1187        };
1188        if let Some((value, version)) = result {
1189            if version <= state.start_version {
1190                state.read_set.insert(key.clone(), version);
1191                return Ok(Some(value));
1192            }
1193            return Ok(None);
1194        }
1195
1196        if let Some(value) = self.manager.sstable_get(key)? {
1197            let start_version = state.start_version;
1198            state.read_set.insert(key.clone(), start_version);
1199            return Ok(Some(value));
1200        }
1201        Ok(None)
1202    }
1203
1204    fn put(&mut self, key: Key, value: Value) -> Result<()> {
1205        let mut state = self
1206            .state
1207            .lock()
1208            .expect("owned memory transaction mutex poisoned");
1209        Self::ensure_active(&state)?;
1210        if state.mode == TxnMode::ReadOnly {
1211            return Err(Error::TxnReadOnly);
1212        }
1213        if state.cursor_open {
1214            return Err(Error::TxnClosed);
1215        }
1216        state.writes.insert(key, Some(value));
1217        Ok(())
1218    }
1219
1220    fn delete(&mut self, key: Key) -> Result<()> {
1221        let mut state = self
1222            .state
1223            .lock()
1224            .expect("owned memory transaction mutex poisoned");
1225        Self::ensure_active(&state)?;
1226        if state.mode == TxnMode::ReadOnly {
1227            return Err(Error::TxnReadOnly);
1228        }
1229        if state.cursor_open {
1230            return Err(Error::TxnClosed);
1231        }
1232        state.writes.insert(key, None);
1233        Ok(())
1234    }
1235
1236    fn scan_prefix(&mut self, prefix: &[u8]) -> Result<Box<dyn OwnedKVScan>> {
1237        // Start at the prefix itself.  An owned cursor advances over the backing map rather
1238        // than a borrowed `BTreeMap::range` guard; starting at keyspace origin would see an
1239        // unrelated earlier key (for example persisted catalog metadata), classify it as out
1240        // of scope, and incorrectly report an empty prefix scan.
1241        self.open_cursor(Some(prefix.to_vec()), Some(prefix.to_vec()), None)
1242    }
1243
1244    fn scan_range(&mut self, start: &[u8], end: &[u8]) -> Result<Box<dyn OwnedKVScan>> {
1245        self.open_cursor(Some(start.to_vec()), None, Some(end.to_vec()))
1246    }
1247
1248    fn commit(self: Box<Self>) -> Result<()> {
1249        let (id, mode, start_version, writes, read_set) = {
1250            let mut state = self
1251                .state
1252                .lock()
1253                .expect("owned memory transaction mutex poisoned");
1254            Self::ensure_active(&state)?;
1255            if state.cursor_open {
1256                return Err(Error::TxnClosed);
1257            }
1258            state.state = TxnState::Committed;
1259            (
1260                state.id,
1261                state.mode,
1262                state.start_version,
1263                std::mem::take(&mut state.writes),
1264                std::mem::take(&mut state.read_set),
1265            )
1266        };
1267        if mode == TxnMode::ReadOnly || writes.is_empty() {
1268            return Ok(());
1269        }
1270
1271        let _snapshot_writer = self.manager.state.owned_snapshot_gate.acquire_writer();
1272        let mut data = self.manager.state.data.write().unwrap();
1273        for key in read_set.keys().chain(writes.keys()) {
1274            let current_version = data.get(key).map(|(_, version)| *version).unwrap_or(0);
1275            if current_version > start_version {
1276                return Err(Error::TxnConflict);
1277            }
1278        }
1279
1280        let mut delta: isize = 0;
1281        for (key, value) in &writes {
1282            let current_size = data
1283                .get(key)
1284                .map(|(value, _)| key.len() + value.len())
1285                .unwrap_or(0);
1286            let new_size = value.as_ref().map_or(0, |value| key.len() + value.len());
1287            delta += new_size as isize - current_size as isize;
1288        }
1289        let current_memory = self.manager.state.current_memory.load(Ordering::Relaxed);
1290        let prospective = if delta >= 0 {
1291            current_memory.saturating_add(delta as usize)
1292        } else {
1293            current_memory.saturating_sub(delta.unsigned_abs())
1294        };
1295        if delta > 0 {
1296            self.manager.state.check_memory_limit(delta as usize)?;
1297        }
1298
1299        let commit_version = self
1300            .manager
1301            .state
1302            .commit_version
1303            .fetch_add(1, Ordering::AcqRel)
1304            + 1;
1305        self.manager.write_wal(id, &writes)?;
1306        for (key, value) in writes {
1307            if let Some(value) = value {
1308                data.insert(key, (value, commit_version));
1309            } else {
1310                data.remove(&key);
1311            }
1312        }
1313        self.manager
1314            .state
1315            .current_memory
1316            .store(prospective, Ordering::Relaxed);
1317        Ok(())
1318    }
1319
1320    fn rollback(self: Box<Self>) -> Result<()> {
1321        let mut state = self
1322            .state
1323            .lock()
1324            .expect("owned memory transaction mutex poisoned");
1325        Self::ensure_active(&state)?;
1326        if state.cursor_open {
1327            return Err(Error::TxnClosed);
1328        }
1329        state.writes.clear();
1330        state.state = TxnState::RolledBack;
1331        Ok(())
1332    }
1333}
1334
1335struct OwnedMemoryCursor {
1336    manager: Arc<MemoryTxnManager>,
1337    transaction: Arc<Mutex<OwnedMemoryTransactionState>>,
1338    snapshot: Option<OwnedSnapshotReader>,
1339    last_key: Option<Key>,
1340    start: Option<Key>,
1341    prefix: Option<Vec<u8>>,
1342    end: Option<Key>,
1343}
1344
1345impl OwnedMemoryCursor {
1346    fn key_is_in_scope(&self, key: &Key) -> bool {
1347        self.prefix
1348            .as_ref()
1349            .is_none_or(|prefix| key.starts_with(prefix))
1350            && self.end.as_ref().is_none_or(|end| key < end)
1351    }
1352
1353    fn data_candidate(&self, start_version: u64) -> Option<(Key, Value, u64)> {
1354        let data = self.manager.state.data.read().unwrap();
1355        let entries: Box<dyn Iterator<Item = (&Key, &(Value, u64))>> = match &self.last_key {
1356            Some(last_key) => Box::new(data.range::<Key, _>((Excluded(last_key), Unbounded))),
1357            None => match &self.start {
1358                Some(start) => Box::new(data.range::<Key, _>((Included(start), Unbounded))),
1359                None => Box::new(data.iter()),
1360            },
1361        };
1362        for (key, (value, version)) in entries {
1363            if !self.key_is_in_scope(key) {
1364                return None;
1365            }
1366            if *version <= start_version {
1367                return Some((key.clone(), value.clone(), *version));
1368            }
1369        }
1370        None
1371    }
1372
1373    fn write_candidate(&self) -> Option<(Key, Option<Value>)> {
1374        let transaction = self
1375            .transaction
1376            .lock()
1377            .expect("owned memory transaction mutex poisoned");
1378        let entry = match &self.last_key {
1379            Some(last_key) => transaction
1380                .writes
1381                .range::<Key, _>((Excluded(last_key), Unbounded))
1382                .next(),
1383            None => match &self.start {
1384                Some(start) => transaction
1385                    .writes
1386                    .range::<Key, _>((Included(start), Unbounded))
1387                    .next(),
1388                None => transaction.writes.iter().next(),
1389            },
1390        }?;
1391        self.key_is_in_scope(entry.0)
1392            .then(|| (entry.0.clone(), entry.1.clone()))
1393    }
1394
1395    fn record_read(&self, key: Key, version: u64) -> Result<()> {
1396        let mut transaction = self
1397            .transaction
1398            .lock()
1399            .expect("owned memory transaction mutex poisoned");
1400        if transaction.state != TxnState::Active || !transaction.cursor_open {
1401            return Err(Error::TxnClosed);
1402        }
1403        transaction.read_set.insert(key, version);
1404        Ok(())
1405    }
1406
1407    fn finish(&mut self) {
1408        self.snapshot.take();
1409        let mut transaction = self
1410            .transaction
1411            .lock()
1412            .expect("owned memory transaction mutex poisoned");
1413        transaction.cursor_open = false;
1414    }
1415}
1416
1417impl OwnedKVScan for OwnedMemoryCursor {
1418    fn next_entry(&mut self) -> Result<Option<(Key, Value)>> {
1419        if self.snapshot.is_none() {
1420            return Ok(None);
1421        }
1422        loop {
1423            let start_version = self
1424                .transaction
1425                .lock()
1426                .expect("owned memory transaction mutex poisoned")
1427                .start_version;
1428            let data = self.data_candidate(start_version);
1429            let write = self.write_candidate();
1430            let next = match (data, write) {
1431                (Some((data_key, data_value, data_version)), Some((write_key, write_value))) => {
1432                    if data_key == write_key {
1433                        self.record_read(data_key.clone(), data_version)?;
1434                        (data_key, write_value)
1435                    } else if data_key < write_key {
1436                        self.record_read(data_key.clone(), data_version)?;
1437                        (data_key, Some(data_value))
1438                    } else {
1439                        (write_key, write_value)
1440                    }
1441                }
1442                (Some((data_key, data_value, data_version)), None) => {
1443                    self.record_read(data_key.clone(), data_version)?;
1444                    (data_key, Some(data_value))
1445                }
1446                (None, Some((write_key, write_value))) => (write_key, write_value),
1447                (None, None) => {
1448                    self.finish();
1449                    return Ok(None);
1450                }
1451            };
1452            self.last_key = Some(next.0.clone());
1453            if let Some(value) = next.1 {
1454                return Ok(Some((next.0, value)));
1455            }
1456        }
1457    }
1458
1459    fn close(&mut self) -> Result<()> {
1460        self.finish();
1461        Ok(())
1462    }
1463}
1464
1465impl Drop for OwnedMemoryCursor {
1466    fn drop(&mut self) {
1467        self.finish();
1468    }
1469}
1470
1471/// Lazy merge iterator that overlays in-flight writes onto a snapshot guard.
1472struct MergedScanIter<'a> {
1473    _data_guard: RwLockReadGuard<'a, BTreeMap<Key, VersionedValue>>,
1474    data_iter: std::collections::btree_map::Range<'a, Key, VersionedValue>,
1475    write_iter: std::collections::btree_map::Range<'a, Key, Option<Value>>,
1476    data_peek: Option<(Key, (Value, u64))>,
1477    write_peek: Option<(Key, Option<Value>)>,
1478    prefix: Option<Vec<u8>>,
1479    end: Option<Key>,
1480    start_version: u64,
1481    read_set: &'a mut HashMap<Key, u64>,
1482}
1483
1484impl<'a> MergedScanIter<'a> {
1485    #[allow(clippy::too_many_arguments)]
1486    fn new(
1487        data_guard: std::sync::RwLockReadGuard<'a, BTreeMap<Key, VersionedValue>>,
1488        data_iter: std::collections::btree_map::Range<'a, Key, VersionedValue>,
1489        write_iter: std::collections::btree_map::Range<'a, Key, Option<Value>>,
1490        prefix: Option<Vec<u8>>,
1491        end: Option<Key>,
1492        start_version: u64,
1493        read_set: &'a mut HashMap<Key, u64>,
1494    ) -> Self {
1495        let mut iter = Self {
1496            _data_guard: data_guard,
1497            data_iter,
1498            write_iter,
1499            data_peek: None,
1500            write_peek: None,
1501            prefix,
1502            end,
1503            start_version,
1504            read_set,
1505        };
1506        iter.advance_data();
1507        iter.advance_write();
1508        iter
1509    }
1510
1511    fn advance_data(&mut self) {
1512        self.data_peek = None;
1513        while let Some((k, (v, ver))) = self.data_iter.next().map(|(k, v)| (k.clone(), v.clone())) {
1514            if let Some(end) = &self.end {
1515                if k >= *end {
1516                    return;
1517                }
1518            }
1519            if let Some(prefix) = &self.prefix {
1520                if !k.starts_with(prefix) {
1521                    return;
1522                }
1523            }
1524            if ver > self.start_version {
1525                continue;
1526            }
1527            self.data_peek = Some((k, (v, ver)));
1528            return;
1529        }
1530    }
1531
1532    fn advance_write(&mut self) {
1533        self.write_peek = None;
1534        if let Some((k, v)) = self.write_iter.next().map(|(k, v)| (k.clone(), v.clone())) {
1535            if let Some(end) = &self.end {
1536                if k >= *end {
1537                    return;
1538                }
1539            }
1540            if let Some(prefix) = &self.prefix {
1541                if !k.starts_with(prefix) {
1542                    return;
1543                }
1544            }
1545            self.write_peek = Some((k, v));
1546        }
1547    }
1548}
1549
1550impl<'a> Iterator for MergedScanIter<'a> {
1551    type Item = (Key, Option<Value>);
1552
1553    fn next(&mut self) -> Option<Self::Item> {
1554        let data_key = self.data_peek.as_ref().map(|(k, _)| k.clone());
1555        let write_key = self.write_peek.as_ref().map(|(k, _)| k.clone());
1556
1557        match (data_key, write_key) {
1558            (Some(dk), Some(wk)) => {
1559                if dk == wk {
1560                    let (_, (_, ver)) = self.data_peek.take().unwrap();
1561                    let (_, write_val) = self.write_peek.take().unwrap();
1562                    self.read_set.insert(dk.clone(), ver);
1563                    self.advance_data();
1564                    self.advance_write();
1565                    Some((dk, write_val))
1566                } else if dk < wk {
1567                    let (k, (v, ver)) = self.data_peek.take().unwrap();
1568                    self.read_set.insert(k.clone(), ver);
1569                    self.advance_data();
1570                    Some((k, Some(v)))
1571                } else {
1572                    let (k, write_val) = self.write_peek.take().unwrap();
1573                    self.advance_write();
1574                    Some((k, write_val))
1575                }
1576            }
1577            (Some(_), None) => {
1578                let (k, (v, ver)) = self.data_peek.take().unwrap();
1579                self.read_set.insert(k.clone(), ver);
1580                self.advance_data();
1581                Some((k, Some(v)))
1582            }
1583            (None, Some(_)) => {
1584                let (k, write_val) = self.write_peek.take().unwrap();
1585                self.advance_write();
1586                Some((k, write_val))
1587            }
1588            (None, None) => None,
1589        }
1590    }
1591}
1592
1593impl<'a> Drop for MemoryTransaction<'a> {
1594    fn drop(&mut self) {
1595        if self.state == TxnState::Active {
1596            self.state = TxnState::RolledBack;
1597        }
1598    }
1599}
1600
1601#[cfg(all(test, not(target_arch = "wasm32")))]
1602mod tests {
1603    use super::*;
1604    use crate::{KVTransaction, TxnManager};
1605    use tempfile::tempdir;
1606    use tracing::Level;
1607
1608    fn key(s: &str) -> Key {
1609        s.as_bytes().to_vec()
1610    }
1611
1612    fn value(s: &str) -> Value {
1613        s.as_bytes().to_vec()
1614    }
1615
1616    fn committed_value_after_reopen(wal_path: &Path, key: Key) -> Option<Value> {
1617        let reopened = MemoryKV::open(wal_path).unwrap();
1618        let manager = reopened.txn_manager();
1619        let mut txn = manager.begin(TxnMode::ReadOnly).unwrap();
1620        txn.get(&key).unwrap()
1621    }
1622
1623    fn write_flush_and_corrupt_sstable<F>(corrupt: F) -> (tempfile::TempDir, PathBuf)
1624    where
1625        F: FnOnce(&Path),
1626    {
1627        let dir = tempdir().unwrap();
1628        let wal_path = dir.path().join("wal.log");
1629        {
1630            let store = MemoryKV::open(&wal_path).unwrap();
1631            let manager = store.txn_manager();
1632            let mut txn = manager.begin(TxnMode::ReadWrite).unwrap();
1633            txn.put(key("k1"), value("v1")).unwrap();
1634            manager.commit(txn).unwrap();
1635            store.flush().unwrap();
1636        }
1637
1638        corrupt(&wal_path.with_extension("sst"));
1639        (dir, wal_path)
1640    }
1641
1642    #[cfg(feature = "test-hooks")]
1643    struct FailsBeforeFsync;
1644
1645    #[cfg(feature = "test-hooks")]
1646    impl IoHooks for FailsBeforeFsync {
1647        fn before_fsync(&self) -> std::io::Result<()> {
1648            Err(std::io::Error::other("injected WAL fsync failure"))
1649        }
1650    }
1651
1652    #[cfg(feature = "test-hooks")]
1653    #[test]
1654    fn commit_self_wal_fsync_failure_does_not_ack_or_apply() {
1655        let dir = tempdir().unwrap();
1656        let wal_path = dir.path().join("wal.log");
1657        let store = MemoryKV::open_with_io_hooks(&wal_path, Arc::new(FailsBeforeFsync)).unwrap();
1658        let manager = store.txn_manager();
1659
1660        let mut txn = manager.begin(TxnMode::ReadWrite).unwrap();
1661        txn.put(key("not-acked"), value("value")).unwrap();
1662        let result = txn.commit_self();
1663        assert!(matches!(result, Err(Error::Io(_))));
1664
1665        let mut read_txn = manager.begin(TxnMode::ReadOnly).unwrap();
1666        assert_eq!(read_txn.get(&key("not-acked")).unwrap(), None);
1667    }
1668
1669    #[test]
1670    fn test_put_and_get_transient() {
1671        let store = MemoryKV::new();
1672        let manager = store.txn_manager();
1673        let mut txn = manager.begin(TxnMode::ReadWrite).unwrap();
1674        txn.put(key("hello"), value("world")).unwrap();
1675        let val = txn.get(&key("hello")).unwrap();
1676        assert_eq!(val, Some(value("world")));
1677        manager.commit(txn).unwrap();
1678
1679        let mut txn2 = manager.begin(TxnMode::ReadOnly).unwrap();
1680        let val2 = txn2.get(&key("hello")).unwrap();
1681        assert_eq!(val2, Some(value("world")));
1682    }
1683
1684    #[test]
1685    fn test_occ_conflict() {
1686        let store = MemoryKV::new();
1687        let manager = store.txn_manager();
1688
1689        let mut t1 = manager.begin(TxnMode::ReadWrite).unwrap();
1690        t1.get(&key("k1")).unwrap();
1691
1692        let mut t2 = manager.begin(TxnMode::ReadWrite).unwrap();
1693        t2.put(key("k1"), value("v2")).unwrap();
1694        assert!(manager.commit(t2).is_ok());
1695
1696        t1.put(key("k1"), value("v1")).unwrap();
1697        let result = manager.commit(t1);
1698        assert!(matches!(result, Err(Error::TxnConflict)));
1699    }
1700
1701    #[test]
1702    fn test_blind_write_conflict() {
1703        let store = MemoryKV::new();
1704        let manager = store.txn_manager();
1705
1706        let mut t1 = manager.begin(TxnMode::ReadWrite).unwrap();
1707        t1.put(key("k1"), value("v1")).unwrap();
1708
1709        let mut t2 = manager.begin(TxnMode::ReadWrite).unwrap();
1710        t2.put(key("k1"), value("v2")).unwrap();
1711        assert!(manager.commit(t2).is_ok());
1712
1713        let result = manager.commit(t1);
1714        assert!(matches!(result, Err(Error::TxnConflict)));
1715    }
1716
1717    #[test]
1718    fn test_read_only_write_fails() {
1719        let store = MemoryKV::new();
1720        let manager = store.txn_manager();
1721        let mut txn = manager.begin(TxnMode::ReadOnly).unwrap();
1722        assert!(matches!(
1723            txn.put(key("k1"), value("v1")),
1724            Err(Error::TxnReadOnly)
1725        ));
1726        assert!(matches!(txn.delete(key("k1")), Err(Error::TxnReadOnly)));
1727    }
1728
1729    #[test]
1730    fn test_txn_closed_error() {
1731        let store = MemoryKV::new();
1732        let manager = store.txn_manager();
1733        let txn = manager.begin(TxnMode::ReadWrite).unwrap();
1734        manager.commit(txn).unwrap();
1735
1736        // This is tricky to test because commit takes ownership.
1737        // We can test by creating a new txn and manually setting its state.
1738        let mut closed_txn = manager.begin(TxnMode::ReadWrite).unwrap();
1739        closed_txn.state = TxnState::Committed;
1740        assert!(matches!(closed_txn.get(&key("k1")), Err(Error::TxnClosed)));
1741        assert!(matches!(
1742            closed_txn.put(key("k1"), value("v1")),
1743            Err(Error::TxnClosed)
1744        ));
1745    }
1746
1747    #[test]
1748    fn test_get_not_found() {
1749        let store = MemoryKV::new();
1750        let manager = store.txn_manager();
1751        let mut txn = manager.begin(TxnMode::ReadOnly).unwrap();
1752        let res = txn.get(&key("non-existent"));
1753        assert!(res.is_ok());
1754        assert!(res.unwrap().is_none());
1755    }
1756
1757    #[test]
1758    fn flush_and_reopen_reads_from_sstable() {
1759        let dir = tempdir().unwrap();
1760        let wal_path = dir.path().join("wal.log");
1761        {
1762            let store = MemoryKV::open(&wal_path).unwrap();
1763            let manager = store.txn_manager();
1764            let mut txn = manager.begin(TxnMode::ReadWrite).unwrap();
1765            txn.put(key("k1"), value("v1")).unwrap();
1766            manager.commit(txn).unwrap();
1767            store.flush().unwrap();
1768        }
1769
1770        let reopened = MemoryKV::open(&wal_path).unwrap();
1771        let manager = reopened.txn_manager();
1772        let mut txn = manager.begin(TxnMode::ReadOnly).unwrap();
1773        assert_eq!(txn.get(&key("k1")).unwrap(), Some(value("v1")));
1774    }
1775
1776    #[test]
1777    fn corrupt_sstable_header_is_discarded_and_wal_recovers() {
1778        let (_dir, wal_path) = write_flush_and_corrupt_sstable(|sst_path| {
1779            let mut file = std::fs::OpenOptions::new()
1780                .write(true)
1781                .open(sst_path)
1782                .unwrap();
1783            use std::io::{Seek, Write};
1784            file.seek(std::io::SeekFrom::Start(0)).unwrap();
1785            file.write_all(b"BAD!").unwrap();
1786            file.sync_all().unwrap();
1787        });
1788        let err = SstableReader::open(&wal_path.with_extension("sst")).unwrap_err();
1789        assert!(matches!(err, Error::InvalidFormat(_)));
1790
1791        assert_eq!(
1792            committed_value_after_reopen(&wal_path, key("k1")),
1793            Some(value("v1"))
1794        );
1795    }
1796
1797    #[test]
1798    fn corrupt_sstable_payload_checksum_is_discarded_and_wal_recovers() {
1799        let (_dir, wal_path) = write_flush_and_corrupt_sstable(|sst_path| {
1800            let mut file = std::fs::OpenOptions::new()
1801                .read(true)
1802                .write(true)
1803                .open(sst_path)
1804                .unwrap();
1805            use std::io::{Read, Seek, Write};
1806            file.seek(std::io::SeekFrom::Start(16 + 8 + key("k1").len() as u64))
1807                .unwrap();
1808            let mut byte = [0u8; 1];
1809            file.read_exact(&mut byte).unwrap();
1810            file.seek(std::io::SeekFrom::Current(-1)).unwrap();
1811            file.write_all(&[byte[0] ^ 0xFF]).unwrap();
1812            file.sync_all().unwrap();
1813        });
1814        let err = SstableReader::open(&wal_path.with_extension("sst")).unwrap_err();
1815        assert!(matches!(err, Error::ChecksumMismatch));
1816
1817        assert_eq!(
1818            committed_value_after_reopen(&wal_path, key("k1")),
1819            Some(value("v1"))
1820        );
1821    }
1822
1823    #[test]
1824    fn truncated_sstable_is_discarded_and_wal_recovers() {
1825        let (_dir, wal_path) = write_flush_and_corrupt_sstable(|sst_path| {
1826            let file = std::fs::OpenOptions::new()
1827                .write(true)
1828                .open(sst_path)
1829                .unwrap();
1830            file.set_len(16).unwrap();
1831            file.sync_all().unwrap();
1832        });
1833        let err = SstableReader::open(&wal_path.with_extension("sst")).unwrap_err();
1834        assert!(matches!(err, Error::InvalidFormat(_)));
1835
1836        assert_eq!(
1837            committed_value_after_reopen(&wal_path, key("k1")),
1838            Some(value("v1"))
1839        );
1840    }
1841
1842    #[test]
1843    fn wal_recovers_committed_tombstone_on_reopen() {
1844        let dir = tempdir().unwrap();
1845        let wal_path = dir.path().join("wal.log");
1846        {
1847            let store = MemoryKV::open(&wal_path).unwrap();
1848            let manager = store.txn_manager();
1849            let mut put_txn = manager.begin(TxnMode::ReadWrite).unwrap();
1850            put_txn.put(key("deleted"), value("value")).unwrap();
1851            manager.commit(put_txn).unwrap();
1852
1853            let mut delete_txn = manager.begin(TxnMode::ReadWrite).unwrap();
1854            delete_txn.delete(key("deleted")).unwrap();
1855            manager.commit(delete_txn).unwrap();
1856        }
1857
1858        assert_eq!(
1859            committed_value_after_reopen(&wal_path, key("deleted")),
1860            None
1861        );
1862    }
1863
1864    #[test]
1865    fn wal_overlays_sstable_on_reopen() {
1866        let dir = tempdir().unwrap();
1867        let wal_path = dir.path().join("wal.log");
1868        {
1869            let store = MemoryKV::open(&wal_path).unwrap();
1870            let manager = store.txn_manager();
1871            let mut txn = manager.begin(TxnMode::ReadWrite).unwrap();
1872            txn.put(key("k1"), value("v1")).unwrap();
1873            manager.commit(txn).unwrap();
1874            store.flush().unwrap();
1875
1876            let mut txn2 = manager.begin(TxnMode::ReadWrite).unwrap();
1877            txn2.put(key("k1"), value("v2")).unwrap();
1878            manager.commit(txn2).unwrap();
1879        }
1880
1881        let reopened = MemoryKV::open(&wal_path).unwrap();
1882        let manager = reopened.txn_manager();
1883        let mut txn = manager.begin(TxnMode::ReadOnly).unwrap();
1884        assert_eq!(txn.get(&key("k1")).unwrap(), Some(value("v2")));
1885    }
1886
1887    #[test]
1888    fn scan_prefix_merges_snapshot_and_writes() {
1889        let store = MemoryKV::new();
1890        let manager = store.txn_manager();
1891
1892        let mut seed = manager.begin(TxnMode::ReadWrite).unwrap();
1893        seed.put(key("p:1"), value("old1")).unwrap();
1894        seed.put(key("p:2"), value("old2")).unwrap();
1895        seed.put(key("q:1"), value("other")).unwrap();
1896        manager.commit(seed).unwrap();
1897
1898        let mut txn = manager.begin(TxnMode::ReadWrite).unwrap();
1899        txn.put(key("p:1"), value("new1")).unwrap();
1900        txn.delete(key("p:2")).unwrap();
1901        txn.put(key("p:3"), value("new3")).unwrap();
1902
1903        let results: Vec<_> = txn.scan_prefix(b"p:").unwrap().collect();
1904        assert_eq!(
1905            results,
1906            vec![(key("p:1"), value("new1")), (key("p:3"), value("new3"))]
1907        );
1908    }
1909
1910    #[test]
1911    fn scan_range_skips_newer_versions() {
1912        let store = MemoryKV::new();
1913        let manager = store.txn_manager();
1914
1915        let mut seed = manager.begin(TxnMode::ReadWrite).unwrap();
1916        seed.put(key("b"), value("v1")).unwrap();
1917        manager.commit(seed).unwrap();
1918
1919        let mut txn1 = manager.begin(TxnMode::ReadWrite).unwrap();
1920
1921        let mut txn2 = manager.begin(TxnMode::ReadWrite).unwrap();
1922        txn2.put(key("ba"), value("v2")).unwrap();
1923        manager.commit(txn2).unwrap();
1924
1925        let results: Vec<_> = txn1.scan_range(b"b", b"c").unwrap().collect();
1926        assert_eq!(results, vec![(key("b"), value("v1"))]);
1927    }
1928
1929    #[test]
1930    fn scan_range_records_reads_for_conflict_detection() {
1931        let store = MemoryKV::new();
1932        let manager = store.txn_manager();
1933
1934        let mut seed = manager.begin(TxnMode::ReadWrite).unwrap();
1935        seed.put(key("k1"), value("v1")).unwrap();
1936        manager.commit(seed).unwrap();
1937
1938        let mut t1 = manager.begin(TxnMode::ReadWrite).unwrap();
1939        let results: Vec<_> = t1.scan_range(b"k0", b"kz").unwrap().collect();
1940        assert_eq!(results, vec![(key("k1"), value("v1"))]);
1941        t1.put(key("k_new"), value("v_new")).unwrap();
1942
1943        let mut t2 = manager.begin(TxnMode::ReadWrite).unwrap();
1944        t2.put(key("k1"), value("v2")).unwrap();
1945        manager.commit(t2).unwrap();
1946
1947        let result = manager.commit(t1);
1948        assert!(matches!(result, Err(Error::TxnConflict)));
1949    }
1950
1951    #[test]
1952    fn owned_memory_transaction_merges_incremental_cursor_and_commits_once() {
1953        use crate::kv::OwnedSessionFactory;
1954
1955        let store = Arc::new(MemoryKV::new());
1956        let session = store
1957            .clone()
1958            .begin_owned_transaction(TxnMode::ReadWrite)
1959            .unwrap();
1960        let lease = session.acquire_lease().unwrap();
1961        lease
1962            .with_transaction(|transaction| {
1963                transaction.put(key("p:1"), value("one"))?;
1964                transaction.put(key("p:2"), value("two"))?;
1965                transaction.put(key("q:1"), value("other"))?;
1966                Ok(())
1967            })
1968            .unwrap();
1969        let mut cursor = lease
1970            .with_transaction(|transaction| transaction.scan_prefix(b"p:"))
1971            .unwrap();
1972        assert_eq!(
1973            cursor.next_entry().unwrap(),
1974            Some((key("p:1"), value("one")))
1975        );
1976        assert_eq!(
1977            cursor.next_entry().unwrap(),
1978            Some((key("p:2"), value("two")))
1979        );
1980        assert_eq!(cursor.next_entry().unwrap(), None);
1981        cursor.close().unwrap();
1982        drop(cursor);
1983        lease
1984            .finish(crate::txn::OwnedLeaseOutcome::Exhausted)
1985            .unwrap();
1986        session.commit().unwrap();
1987
1988        let read = store
1989            .clone()
1990            .begin_owned_read(crate::kv::OwnedReadOptions::default())
1991            .unwrap();
1992        let lease = read.acquire_lease().unwrap();
1993        assert_eq!(
1994            lease
1995                .with_transaction(|transaction| transaction.get(&key("p:2")))
1996                .unwrap(),
1997            Some(value("two"))
1998        );
1999        lease
2000            .finish(crate::txn::OwnedLeaseOutcome::Exhausted)
2001            .unwrap();
2002
2003        let range = store
2004            .clone()
2005            .begin_owned_read(crate::kv::OwnedReadOptions::default())
2006            .unwrap();
2007        let lease = range.acquire_lease().unwrap();
2008        let mut cursor = lease
2009            .with_transaction(|transaction| transaction.scan_range(b"p:1", b"q:"))
2010            .unwrap();
2011        assert_eq!(
2012            cursor.next_entry().unwrap(),
2013            Some((key("p:1"), value("one")))
2014        );
2015        assert_eq!(
2016            cursor.next_entry().unwrap(),
2017            Some((key("p:2"), value("two")))
2018        );
2019        assert_eq!(cursor.next_entry().unwrap(), None);
2020        drop(cursor);
2021        lease
2022            .finish(crate::txn::OwnedLeaseOutcome::Exhausted)
2023            .unwrap();
2024    }
2025
2026    #[test]
2027    fn owned_prefix_cursor_skips_keys_before_its_prefix() {
2028        use crate::kv::OwnedSessionFactory;
2029
2030        let store = Arc::new(MemoryKV::new());
2031        let writer = store
2032            .clone()
2033            .begin_owned_transaction(TxnMode::ReadWrite)
2034            .unwrap();
2035        let writer_lease = writer.acquire_lease().unwrap();
2036        writer_lease
2037            .with_transaction(|transaction| {
2038                transaction.put(key("catalog:before"), value("metadata"))?;
2039                transaction.put(key("row:1"), value("one"))?;
2040                Ok(())
2041            })
2042            .unwrap();
2043        writer_lease
2044            .finish(crate::txn::OwnedLeaseOutcome::Exhausted)
2045            .unwrap();
2046        writer.commit().unwrap();
2047
2048        let reader = store.clone().begin_owned_read(Default::default()).unwrap();
2049        let reader_lease = reader.acquire_lease().unwrap();
2050        let mut cursor = reader_lease
2051            .with_transaction(|transaction| transaction.scan_prefix(b"row:"))
2052            .unwrap();
2053        assert_eq!(
2054            cursor.next_entry().unwrap(),
2055            Some((key("row:1"), value("one")))
2056        );
2057        assert_eq!(cursor.next_entry().unwrap(), None);
2058        cursor.close().unwrap();
2059        drop(cursor);
2060        reader_lease
2061            .finish(crate::txn::OwnedLeaseOutcome::Exhausted)
2062            .unwrap();
2063    }
2064
2065    #[test]
2066    fn any_kv_memory_dispatches_to_owned_memory_session_without_borrowed_transaction() {
2067        use crate::kv::{AnyKV, OwnedSessionFactory};
2068
2069        let store = Arc::new(AnyKV::Memory(MemoryKV::new()));
2070        let session = store
2071            .clone()
2072            .begin_owned_transaction(TxnMode::ReadWrite)
2073            .unwrap();
2074        let lease = session.acquire_lease().unwrap();
2075        lease
2076            .with_transaction(|transaction| transaction.put(key("owned"), value("value")))
2077            .unwrap();
2078        lease
2079            .finish(crate::txn::OwnedLeaseOutcome::Exhausted)
2080            .unwrap();
2081        session.commit().unwrap();
2082
2083        let read = store
2084            .clone()
2085            .begin_owned_read(crate::kv::OwnedReadOptions::default())
2086            .unwrap();
2087        let lease = read.acquire_lease().unwrap();
2088        assert_eq!(
2089            lease
2090                .with_transaction(|transaction| transaction.get(&key("owned")))
2091                .unwrap(),
2092            Some(value("value"))
2093        );
2094        lease
2095            .finish(crate::txn::OwnedLeaseOutcome::Exhausted)
2096            .unwrap();
2097    }
2098
2099    #[test]
2100    fn memory_stats_tracks_put_and_delete() {
2101        let store = MemoryKV::new();
2102        let manager = store.txn_manager();
2103
2104        let stats = manager.memory_stats();
2105        assert_eq!(stats.total_bytes, 0);
2106        assert_eq!(stats.kv_bytes, 0);
2107        assert_eq!(stats.index_bytes, 0);
2108
2109        // Insert a value and commit.
2110        let mut txn = manager.begin(TxnMode::ReadWrite).unwrap();
2111        txn.put(key("a"), value("1234")).unwrap(); // key=1, value=4 => 5 bytes
2112        manager.commit(txn).unwrap();
2113
2114        let stats = manager.memory_stats();
2115        assert_eq!(stats.total_bytes, 5);
2116        assert_eq!(stats.kv_bytes, 5);
2117        assert_eq!(stats.index_bytes, 0);
2118
2119        // Delete and ensure usage returns to zero.
2120        let mut txn = manager.begin(TxnMode::ReadWrite).unwrap();
2121        txn.delete(key("a")).unwrap();
2122        manager.commit(txn).unwrap();
2123
2124        let stats = manager.memory_stats();
2125        assert_eq!(stats.total_bytes, 0);
2126        assert_eq!(stats.kv_bytes, 0);
2127    }
2128
2129    #[test]
2130    fn memory_limit_error_does_not_break_reads() {
2131        let store = MemoryKV::new_with_limit(Some(10));
2132        let manager = store.txn_manager();
2133
2134        // First insert within limit: key(2) + value(4) = 6.
2135        let mut txn = manager.begin_internal(TxnMode::ReadWrite).unwrap();
2136        txn.put(key("k1"), value("vvvv")).unwrap();
2137        manager.commit(txn).unwrap();
2138
2139        // Next insert would exceed limit: key(2) + value(6) + existing(6) -> 14 > 10.
2140        let mut txn2 = manager.begin_internal(TxnMode::ReadWrite).unwrap();
2141        txn2.put(key("k2"), value("vvvvvv")).unwrap();
2142        let result = manager.commit(txn2);
2143        assert!(matches!(result, Err(Error::MemoryLimitExceeded { .. })));
2144
2145        // Read still works and existing data intact.
2146        let mut read_txn = manager.begin_internal(TxnMode::ReadOnly).unwrap();
2147        let got = read_txn.get(&key("k1")).unwrap();
2148        assert_eq!(got, Some(value("vvvv")));
2149
2150        // Memory usage stays at the previous successful commit.
2151        let stats = manager.memory_stats();
2152        assert_eq!(stats.total_bytes, 6);
2153    }
2154
2155    struct VecWriter(std::sync::Arc<std::sync::Mutex<Vec<u8>>>);
2156
2157    impl std::io::Write for VecWriter {
2158        fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
2159            let mut guard = self.0.lock().unwrap();
2160            guard.extend_from_slice(buf);
2161            Ok(buf.len())
2162        }
2163
2164        fn flush(&mut self) -> std::io::Result<()> {
2165            Ok(())
2166        }
2167    }
2168
2169    #[test]
2170    fn compaction_skips_when_over_limit_and_logs_warning() {
2171        let store = MemoryKV::new_with_limit(Some(12));
2172        let manager = store.txn_manager();
2173
2174        // Populate data to track current memory: key(2)+val(6)=8 bytes.
2175        let mut txn = manager.begin_internal(TxnMode::ReadWrite).unwrap();
2176        txn.put(key("k1"), value("123456")).unwrap();
2177        manager.commit(txn).unwrap();
2178
2179        // Prepare log capture.
2180        let buffer = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2181        let make_writer = {
2182            let buf = buffer.clone();
2183            move || VecWriter(buf.clone())
2184        };
2185        let subscriber = tracing_subscriber::fmt()
2186            .with_max_level(Level::WARN)
2187            .with_writer(make_writer)
2188            .without_time()
2189            .finish();
2190        let _guard = tracing::subscriber::set_default(subscriber);
2191
2192        // input=2 (assume one entry), output=10 => projected 8-2+10=16 > 12 -> skip.
2193        let ran = manager.compact_with_limit(2, 10, || Ok(())).unwrap();
2194        assert!(!ran);
2195
2196        // Memory usage unchanged.
2197        assert_eq!(manager.memory_stats().total_bytes, 8);
2198
2199        // Verify warning was logged.
2200        let log = String::from_utf8(buffer.lock().unwrap().clone()).unwrap();
2201        assert!(
2202            log.contains("compaction skipped due to memory limit"),
2203            "expected warning log, got: {}",
2204            log
2205        );
2206    }
2207}