1use 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#[derive(Debug, Clone, Default, PartialEq)]
22pub struct MemoryStats {
23 pub total_bytes: usize,
25 pub kv_bytes: usize,
27 pub index_bytes: usize,
29}
30
31#[derive(Clone)]
33pub struct MemoryKV {
34 manager: Arc<MemoryTxnManager>,
35}
36
37impl MemoryKV {
38 pub fn new() -> Self {
40 Self {
41 manager: Arc::new(MemoryTxnManager::new(None, None, None)),
42 }
43 }
44
45 pub fn memory_stats(&self) -> MemoryStats {
47 self.manager.memory_stats()
48 }
49
50 pub fn new_with_limit(limit: Option<usize>) -> Self {
52 Self {
53 manager: Arc::new(MemoryTxnManager::new_with_limit(limit)),
54 }
55 }
56
57 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 #[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 #[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 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
146type VersionedValue = (Value, u64);
148
149struct MemorySharedState {
151 data: RwLock<BTreeMap<Key, VersionedValue>>,
153 next_txn_id: AtomicU64,
155 commit_version: AtomicU64,
157 wal_writer: Option<RwLock<WalWriter>>,
159 wal_path: Option<PathBuf>,
161 sstable: RwLock<Option<SstableReader>>,
163 sstable_path: Option<PathBuf>,
165 memory_limit: RwLock<Option<usize>>,
167 current_memory: AtomicUsize,
169 owned_snapshot_gate: Arc<OwnedSnapshotGate>,
172 #[cfg(feature = "test-hooks")]
173 io_hooks: RwLock<Option<Arc<dyn IoHooks>>>,
175 #[cfg(feature = "test-hooks")]
176 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 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 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 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
304pub 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 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 pub fn memory_stats(&self) -> MemoryStats {
373 self.state.memory_stats()
374 }
375
376 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 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 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 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 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 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 wal.sync()?;
520 }
521 Ok(())
522 }
523
524 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 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 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 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 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 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 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 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
835pub 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 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 (&*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 (&*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 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 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 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 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 self.manager.write_wal(self.id, &self.writes)?;
1062
1063 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
1090struct 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 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
1471struct 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 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 let mut txn = manager.begin(TxnMode::ReadWrite).unwrap();
2111 txn.put(key("a"), value("1234")).unwrap(); 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 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 let mut txn = manager.begin_internal(TxnMode::ReadWrite).unwrap();
2136 txn.put(key("k1"), value("vvvv")).unwrap();
2137 manager.commit(txn).unwrap();
2138
2139 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 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 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 let mut txn = manager.begin_internal(TxnMode::ReadWrite).unwrap();
2176 txn.put(key("k1"), value("123456")).unwrap();
2177 manager.commit(txn).unwrap();
2178
2179 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 let ran = manager.compact_with_limit(2, 10, || Ok(())).unwrap();
2194 assert!(!ran);
2195
2196 assert_eq!(manager.memory_stats().total_bytes, 8);
2198
2199 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}