Skip to main content

mentedb_storage/
engine.rs

1//! Storage Engine: facade that ties the page manager, WAL, and buffer pool together.
2
3use std::fs::File;
4use std::path::Path;
5
6use mentedb_core::MemoryNode;
7use mentedb_core::error::{MenteError, MenteResult};
8
9use parking_lot::Mutex;
10use tracing::info;
11
12use crate::buffer::BufferPool;
13use crate::page::{PAGE_DATA_SIZE, Page, PageId, PageManager, PageType};
14use crate::wal::{Wal, WalEntryType};
15/// Default number of page frames in the buffer pool.
16const DEFAULT_BUFFER_POOL_SIZE: usize = 1024;
17
18/// Auto-checkpoint when WAL file exceeds this size (8 MB).
19const WAL_AUTO_CHECKPOINT_BYTES: u64 = 8 * 1024 * 1024;
20
21/// The unified storage engine for MenteDB.
22///
23/// Coordinates page allocation, caching, and write-ahead logging to provide
24/// crash-safe, page-oriented storage for memory nodes.
25///
26/// Concurrency model (inspired by WAL-mode databases):
27/// - **Reads are lock-free**: `read_page` only touches the buffer pool and page
28///   manager — no file locks, no WAL access.
29/// - **Writes are fully serialized** via a blocking `flock(2)` on the WAL file.
30///   The entire write transaction (page allocation + WAL append + page write +
31///   fsync) executes under a single flock, ensuring correctness across multiple
32///   processes sharing the same data directory.
33/// - **State is refreshed from disk** under the flock: page count is re-read
34///   from the file header and LSN is re-read from the WAL tail, so no process
35///   can act on stale in-memory state.
36/// - **One process per database.** Open acquires an exclusive lock on a LOCK
37///   file in the data directory and holds it for the engine's lifetime. The
38///   per operation flock protects individual write transactions, but layers
39///   above the storage engine cache state across operations (buffer pool,
40///   page map, index snapshots), and a second concurrent process flushing
41///   those caches rolls the database back to its own open time snapshot.
42///   A second open therefore fails fast with a "locked" error instead of
43///   silently corrupting state; callers retry until the owner exits.
44pub struct StorageEngine {
45    page_manager: Mutex<PageManager>,
46    buffer_pool: BufferPool,
47    wal: Mutex<Wal>,
48    /// Held exclusively for the lifetime of this engine; released on close
49    /// or process exit. Guards against concurrent multi process opens.
50    process_lock: Mutex<Option<File>>,
51}
52
53impl StorageEngine {
54    /// Open (or create) a storage engine rooted at `path`.
55    ///
56    /// `path` must be a directory; it will be created if it does not exist.
57    /// After opening, any uncommitted WAL entries are replayed for crash recovery.
58    ///
59    /// # Example
60    ///
61    /// ```no_run
62    /// use mentedb_storage::StorageEngine;
63    ///
64    /// let engine = StorageEngine::open("/tmp/mentedb".as_ref())?;
65    /// // engine is ready — WAL recovery already ran if needed
66    /// # Ok::<(), mentedb_core::error::MenteError>(())
67    /// ```
68    pub fn open(path: &Path) -> MenteResult<Self> {
69        std::fs::create_dir_all(path)?;
70
71        // Exclusive process lock, held until close or process exit. A second
72        // process opening the same directory would interleave stale cached
73        // state with ours and roll the database back, so refuse it up front.
74        let lock_path = path.join("LOCK");
75        let lock_file = std::fs::OpenOptions::new()
76            .create(true)
77            .truncate(false)
78            .write(true)
79            .open(&lock_path)?;
80        // Cross-host lock: fcntl OFD lock on Linux (enforced across hosts by the
81        // NFSv4 lock manager), flock elsewhere. This makes single-writer safe
82        // when the directory lives on a network filesystem mounted by more than
83        // one host.
84        match crate::lock::try_lock_exclusive(&lock_file) {
85            Ok(true) => {}
86            Ok(false) => {
87                return Err(MenteError::Storage(format!(
88                    "database directory {} is locked by another process",
89                    path.display()
90                )));
91            }
92            Err(e) => {
93                return Err(MenteError::Storage(format!(
94                    "failed to lock database directory {}: {e}",
95                    path.display()
96                )));
97            }
98        }
99
100        let page_manager = PageManager::open(path)?;
101        let buffer_pool = BufferPool::new(DEFAULT_BUFFER_POOL_SIZE);
102        let wal = Wal::open(path)?;
103
104        let engine = Self {
105            page_manager: Mutex::new(page_manager),
106            buffer_pool,
107            wal: Mutex::new(wal),
108            process_lock: Mutex::new(Some(lock_file)),
109        };
110
111        let recovered = engine.recover()?;
112        if recovered > 0 {
113            info!(recovered, ?path, "storage engine opened with WAL recovery");
114        } else {
115            info!(?path, "storage engine opened");
116        }
117
118        Ok(engine)
119    }
120
121    /// Replay WAL entries to recover writes that were not checkpointed.
122    ///
123    /// For each `PageWrite` entry the serialized data is written back to its page.
124    /// After replay the WAL is truncated. Returns the number of entries replayed.
125    pub fn recover(&self) -> MenteResult<usize> {
126        let mut wal = self.wal.lock();
127        wal.lock_exclusive()?;
128        let entries = wal.iterate()?;
129        let mut count = 0usize;
130        let mut pm = self.page_manager.lock();
131
132        // Refresh page count from disk — another process may have written pages.
133        pm.reload_header()?;
134
135        // Replay with last-op-wins per page: every PageWrite carries a full page
136        // image and PageFree discards the page, so only the final entry for each
137        // page matters. This is what makes free-then-reuse sequences safe: a
138        // PageFree followed by a later PageWrite for the same page must not
139        // leave the page on the free list.
140        let mut last_op: std::collections::HashMap<u64, &crate::wal::WalEntry> = Default::default();
141        let mut order: Vec<u64> = Vec::new();
142        for entry in &entries {
143            match entry.entry_type {
144                WalEntryType::PageWrite | WalEntryType::PageFree => {
145                    if !last_op.contains_key(&entry.page_id) {
146                        order.push(entry.page_id);
147                    }
148                    last_op.insert(entry.page_id, entry);
149                }
150                WalEntryType::Checkpoint | WalEntryType::Commit => {}
151            }
152        }
153
154        for page_id_raw in order {
155            let entry = last_op[&page_id_raw];
156            let page_id = PageId(entry.page_id);
157            match entry.entry_type {
158                WalEntryType::PageWrite => {
159                    while pm.page_count() <= entry.page_id {
160                        pm.allocate_page()?;
161                    }
162
163                    let mut page = pm.read_page(page_id)?;
164                    let copy_len = entry.data.len().min(PAGE_DATA_SIZE);
165                    page.data[..copy_len].copy_from_slice(&entry.data[..copy_len]);
166                    if copy_len < PAGE_DATA_SIZE {
167                        page.data[copy_len..].fill(0);
168                    }
169                    page.header.page_id = entry.page_id;
170                    page.header.lsn = entry.lsn;
171                    page.header.page_type = PageType::Data as u8;
172                    page.header.free_space = (PAGE_DATA_SIZE - copy_len) as u16;
173                    page.header.checksum = page.compute_checksum();
174
175                    pm.write_page(page_id, &page)?;
176                    count += 1;
177                }
178                WalEntryType::PageFree => {
179                    // Mark the page Free without touching the free list; the
180                    // list is rebuilt from page types after replay so it stays
181                    // consistent with WAL-derived page states.
182                    if entry.page_id < pm.page_count() {
183                        let mut page = Page::zeroed();
184                        page.header.page_id = entry.page_id;
185                        page.header.page_type = PageType::Free as u8;
186                        pm.write_page(page_id, &page)?;
187                        self.buffer_pool.invalidate(page_id);
188                        count += 1;
189                    }
190                }
191                WalEntryType::Checkpoint | WalEntryType::Commit => {}
192            }
193        }
194
195        if count > 0 {
196            // Header updates (free list head, page count) are not fsynced on
197            // every write, so after a crash the free list can disagree with the
198            // replayed page states. Rebuild it from page types: the WAL is the
199            // sole source of truth for which pages are Free vs Data.
200            pm.rebuild_free_list()?;
201            pm.sync()?;
202            let next_lsn = wal.next_lsn();
203            wal.truncate(next_lsn)?;
204            info!(count, "WAL recovery replayed entries");
205        }
206
207        wal.unlock()?;
208        Ok(count)
209    }
210
211    /// Gracefully shut down: flush dirty pages, sync files.
212    ///
213    /// # Example
214    ///
215    /// ```no_run
216    /// # use mentedb_storage::StorageEngine;
217    /// # let engine = StorageEngine::open("/tmp/mentedb".as_ref())?;
218    /// engine.close()?;
219    /// # Ok::<(), mentedb_core::error::MenteError>(())
220    /// ```
221    pub fn close(&self) -> MenteResult<()> {
222        let mut pm = self.page_manager.lock();
223        self.buffer_pool.flush_all(&mut pm)?;
224        pm.sync()?;
225        self.wal.lock().sync()?;
226        // Release the process lock so the directory can be reopened, by this
227        // process or another, without waiting for us to exit.
228        if let Some(lock_file) = self.process_lock.lock().take() {
229            let _ = crate::lock::unlock(&lock_file);
230        }
231        info!("storage engine closed");
232        Ok(())
233    }
234
235    /// Release the process lock without flushing anything.
236    ///
237    /// Test support for simulating a process crash: a real crash releases
238    /// the OS file lock but persists nothing beyond what was already synced.
239    #[doc(hidden)]
240    pub fn release_process_lock(&self) {
241        if let Some(lock_file) = self.process_lock.lock().take() {
242            let _ = crate::lock::unlock(&lock_file);
243        }
244    }
245
246    // ---- low-level page operations ----
247
248    /// Allocate a fresh page (for internal/test use).
249    ///
250    /// **WARNING**: In multi-process scenarios, prefer `store_memory` which
251    /// allocates under the WAL flock. This method does NOT acquire the flock.
252    pub fn allocate_page(&self) -> MenteResult<PageId> {
253        self.page_manager.lock().allocate_page()
254    }
255
256    /// Read a page through the buffer pool (lock-free — no WAL access).
257    pub fn read_page(&self, page_id: PageId) -> MenteResult<Box<Page>> {
258        self.buffer_pool
259            .fetch_page(page_id, &mut self.page_manager.lock())
260    }
261
262    /// Buffer-pool cache activity, for metrics (hit ratio, evictions, residency).
263    pub fn buffer_stats(&self) -> crate::buffer::BufferStats {
264        self.buffer_pool.stats()
265    }
266
267    /// Total pages in the store; multiply by `PAGE_SIZE` for on-disk bytes.
268    pub fn page_count(&self) -> u64 {
269        self.page_manager.lock().page_count()
270    }
271
272    /// Write data into an already-allocated page with WAL protection.
273    ///
274    /// Acquires the WAL flock for the duration of the write transaction.
275    /// For new pages, prefer `store_memory` which allocates + writes atomically.
276    pub fn write_page(&self, page_id: PageId, data: &[u8]) -> MenteResult<()> {
277        let lsn = {
278            let mut wal = self.wal.lock();
279            wal.lock_exclusive()?;
280            wal.reload_lsn()?;
281            let lsn = wal.append(WalEntryType::PageWrite, page_id.0, data)?;
282            wal.sync()?;
283            wal.unlock()?;
284            lsn
285        };
286
287        self.apply_page_write(page_id, data, lsn)
288    }
289
290    /// Apply a page write to the buffer pool and page manager (after WAL).
291    fn apply_page_write(&self, page_id: PageId, data: &[u8], lsn: u64) -> MenteResult<()> {
292        let mut pm = self.page_manager.lock();
293        let mut page = self.buffer_pool.fetch_page(page_id, &mut pm)?;
294        drop(pm);
295
296        let copy_len = data.len().min(PAGE_DATA_SIZE);
297        page.data[..copy_len].copy_from_slice(&data[..copy_len]);
298        if copy_len < PAGE_DATA_SIZE {
299            page.data[copy_len..].fill(0);
300        }
301        page.header.lsn = lsn;
302        page.header.page_type = PageType::Data as u8;
303        page.header.free_space = (PAGE_DATA_SIZE - copy_len) as u16;
304        page.header.checksum = page.compute_checksum();
305
306        if self.buffer_pool.update_page(page_id, &page).is_err() {
307            self.page_manager.lock().write_page(page_id, &page)?;
308        }
309        self.buffer_pool.unpin_page(page_id, true).ok();
310
311        Ok(())
312    }
313
314    // ---- high-level memory operations ----
315
316    /// Serialize and store a [`MemoryNode`] into a single page.
317    ///
318    /// The entire operation — page allocation, WAL append, page write — executes
319    /// under a single WAL flock, making it safe across multiple processes.
320    ///
321    /// # Example
322    ///
323    /// ```no_run
324    /// use mentedb_storage::StorageEngine;
325    /// use mentedb_core::{MemoryNode, memory::MemoryType, types::AgentId};
326    ///
327    /// let engine = StorageEngine::open("/tmp/mentedb".as_ref())?;
328    /// let node = MemoryNode::new(
329    ///     AgentId::new(),
330    ///     MemoryType::Semantic,
331    ///     "User likes dark mode".to_string(),
332    ///     vec![0.1, 0.2],
333    /// );
334    /// let page_id = engine.store_memory(&node)?;
335    /// # Ok::<(), mentedb_core::error::MenteError>(())
336    /// ```
337    pub fn store_memory(&self, node: &MemoryNode) -> MenteResult<PageId> {
338        let serialized =
339            serde_json::to_vec(node).map_err(|e| MenteError::Serialization(e.to_string()))?;
340
341        if serialized.len() + 4 > PAGE_DATA_SIZE {
342            return Err(MenteError::CapacityExceeded(format!(
343                "memory node serialized to {} bytes (max {})",
344                serialized.len(),
345                PAGE_DATA_SIZE - 4,
346            )));
347        }
348
349        let mut buf = Vec::with_capacity(4 + serialized.len());
350        buf.extend_from_slice(&(serialized.len() as u32).to_le_bytes());
351        buf.extend_from_slice(&serialized);
352
353        // Atomic write transaction: allocate + WAL + page write under one flock
354        let (page_id, lsn) = {
355            let mut wal = self.wal.lock();
356            let mut pm = self.page_manager.lock();
357
358            // Acquire flock and refresh state from disk
359            wal.lock_exclusive()?;
360            pm.reload_header()?;
361            wal.reload_lsn()?;
362
363            // Allocate page (using fresh page_count from disk)
364            let page_id = pm.allocate_page()?;
365
366            // WAL append + sync (WAL fsync guarantees durability;
367            // page data is written but not fsynced — checkpoint handles that)
368            let lsn = wal.append(WalEntryType::PageWrite, page_id.0, &buf)?;
369            wal.sync()?;
370
371            // Write page data to disk (no fsync — recoverable from WAL)
372            let mut page = Page::zeroed();
373            page.header.page_id = page_id.0;
374            let copy_len = buf.len().min(PAGE_DATA_SIZE);
375            page.data[..copy_len].copy_from_slice(&buf[..copy_len]);
376            page.header.lsn = lsn;
377            page.header.page_type = PageType::Data as u8;
378            page.header.free_space = (PAGE_DATA_SIZE - copy_len) as u16;
379            page.header.checksum = page.compute_checksum();
380            pm.write_page(page_id, &page)?;
381
382            // Release flock — other processes can now write
383            wal.unlock()?;
384
385            (page_id, lsn)
386        };
387
388        // Drop any stale cached copy (the page may have been reused from the
389        // free list); the next read loads the fresh data from disk.
390        let _ = lsn;
391        self.buffer_pool.invalidate(page_id);
392
393        // Auto-checkpoint when WAL exceeds threshold to prevent unbounded growth.
394        // This keeps reload_lsn() fast for subsequent writes.
395        if self.wal.lock().file_size() > WAL_AUTO_CHECKPOINT_BYTES
396            && let Err(e) = self.checkpoint()
397        {
398            tracing::warn!("auto-checkpoint failed: {e}");
399        }
400
401        info!(
402            page_id = page_id.0,
403            bytes = serialized.len(),
404            "stored memory node"
405        );
406        Ok(page_id)
407    }
408
409    /// Store multiple [`MemoryNode`]s in a single locked transaction.
410    ///
411    /// Acquires the WAL flock once, writes all nodes, then releases. This avoids
412    /// the per-write overhead of `reload_header` / `reload_lsn` for bulk inserts.
413    /// Auto-checkpoints after the batch if the WAL exceeds the threshold.
414    pub fn store_memory_batch(&self, nodes: &[MemoryNode]) -> MenteResult<Vec<PageId>> {
415        // Phase 1: serialize all nodes upfront (no locks held)
416        let mut bufs = Vec::with_capacity(nodes.len());
417        for node in nodes {
418            let serialized =
419                serde_json::to_vec(node).map_err(|e| MenteError::Serialization(e.to_string()))?;
420            if serialized.len() + 4 > PAGE_DATA_SIZE {
421                return Err(MenteError::CapacityExceeded(format!(
422                    "memory node serialized to {} bytes (max {})",
423                    serialized.len(),
424                    PAGE_DATA_SIZE - 4,
425                )));
426            }
427            let mut buf = Vec::with_capacity(4 + serialized.len());
428            buf.extend_from_slice(&(serialized.len() as u32).to_le_bytes());
429            buf.extend_from_slice(&serialized);
430            bufs.push(buf);
431        }
432
433        // Phase 2: single locked transaction for all writes
434        let page_ids = {
435            let mut wal = self.wal.lock();
436            let mut pm = self.page_manager.lock();
437
438            wal.lock_exclusive()?;
439            pm.reload_header()?;
440            wal.reload_lsn()?;
441
442            let mut ids = Vec::with_capacity(bufs.len());
443            for buf in &bufs {
444                let page_id = pm.allocate_page()?;
445                let lsn = wal.append(WalEntryType::PageWrite, page_id.0, buf)?;
446
447                let mut page = Page::zeroed();
448                page.header.page_id = page_id.0;
449                let copy_len = buf.len().min(PAGE_DATA_SIZE);
450                page.data[..copy_len].copy_from_slice(&buf[..copy_len]);
451                page.header.lsn = lsn;
452                page.header.page_type = PageType::Data as u8;
453                page.header.free_space = (PAGE_DATA_SIZE - copy_len) as u16;
454                page.header.checksum = page.compute_checksum();
455                pm.write_page(page_id, &page)?;
456
457                ids.push(page_id);
458            }
459
460            // WAL fsync only — page data is recoverable from WAL on crash.
461            // Checkpoint handles page file fsync.
462            wal.sync()?;
463            wal.unlock()?;
464
465            ids
466        };
467
468        // Drop stale cached copies for pages reused from the free list.
469        for page_id in &page_ids {
470            self.buffer_pool.invalidate(*page_id);
471        }
472
473        // Auto-checkpoint if WAL grew too large
474        if self.wal.lock().file_size() > WAL_AUTO_CHECKPOINT_BYTES
475            && let Err(e) = self.checkpoint()
476        {
477            tracing::warn!("auto-checkpoint failed: {e}");
478        }
479
480        info!(count = page_ids.len(), "stored memory batch");
481        Ok(page_ids)
482    }
483
484    /// Update a [`MemoryNode`] in place on its existing page.
485    ///
486    /// The write goes through the WAL-protected `write_page` path, so it is
487    /// crash-durable and keeps the buffer pool coherent. Unlike storing to a
488    /// fresh page, this never orphans the old copy.
489    pub fn update_memory(&self, page_id: PageId, node: &MemoryNode) -> MenteResult<()> {
490        let serialized =
491            serde_json::to_vec(node).map_err(|e| MenteError::Serialization(e.to_string()))?;
492
493        if serialized.len() + 4 > PAGE_DATA_SIZE {
494            return Err(MenteError::CapacityExceeded(format!(
495                "memory node serialized to {} bytes (max {})",
496                serialized.len(),
497                PAGE_DATA_SIZE - 4,
498            )));
499        }
500
501        let mut buf = Vec::with_capacity(4 + serialized.len());
502        buf.extend_from_slice(&(serialized.len() as u32).to_le_bytes());
503        buf.extend_from_slice(&serialized);
504
505        self.write_page(page_id, &buf)
506    }
507
508    /// Delete a memory by returning its page to the free list.
509    ///
510    /// The deletion is WAL-logged before the page is freed, so it survives a
511    /// crash: recovery replays the `PageFree` entry and the memory does not
512    /// resurrect on reopen. The freed page is reused by later allocations.
513    pub fn delete_memory(&self, page_id: PageId) -> MenteResult<()> {
514        {
515            let mut wal = self.wal.lock();
516            let mut pm = self.page_manager.lock();
517
518            wal.lock_exclusive()?;
519            pm.reload_header()?;
520            wal.reload_lsn()?;
521
522            // WAL fsync guarantees the deletion is durable before the page
523            // is touched; the free-list update itself is recoverable.
524            wal.append(WalEntryType::PageFree, page_id.0, &[])?;
525            wal.sync()?;
526
527            pm.free_page(page_id)?;
528            wal.unlock()?;
529        }
530
531        // A stale cached copy must never be served or flushed back.
532        self.buffer_pool.invalidate(page_id);
533
534        info!(page_id = page_id.0, "deleted memory node");
535        Ok(())
536    }
537
538    /// Load and deserialize a [`MemoryNode`] from the given page.
539    ///
540    /// # Example
541    ///
542    /// ```no_run
543    /// # use mentedb_storage::{StorageEngine, PageId};
544    /// # let engine = StorageEngine::open("/tmp/mentedb".as_ref())?;
545    /// let node = engine.load_memory(PageId(1))?;
546    /// println!("memory: {}", node.content);
547    /// # Ok::<(), mentedb_core::error::MenteError>(())
548    /// ```
549    pub fn load_memory(&self, page_id: PageId) -> MenteResult<MemoryNode> {
550        let page = self.read_page(page_id)?;
551        self.buffer_pool.unpin_page(page_id, false).ok();
552
553        if PageType::from(page.header.page_type) != PageType::Data {
554            return Err(MenteError::Storage(format!(
555                "page {} is not a data page",
556                page_id.0
557            )));
558        }
559
560        let len = u32::from_le_bytes(page.data[..4].try_into().unwrap()) as usize;
561        if len == 0 || len + 4 > PAGE_DATA_SIZE {
562            return Err(MenteError::Storage(format!(
563                "invalid memory node length prefix: {len}"
564            )));
565        }
566
567        serde_json::from_slice(&page.data[4..4 + len])
568            .map_err(|e| MenteError::Serialization(e.to_string()))
569    }
570
571    // ---- durability ----
572
573    /// Checkpoint: flush all dirty pages, sync to disk, and truncate the WAL.
574    ///
575    /// # Example
576    ///
577    /// ```no_run
578    /// # use mentedb_storage::StorageEngine;
579    /// # let engine = StorageEngine::open("/tmp/mentedb".as_ref())?;
580    /// // After a batch of writes, checkpoint to reclaim WAL space
581    /// engine.checkpoint()?;
582    /// # Ok::<(), mentedb_core::error::MenteError>(())
583    /// ```
584    pub fn checkpoint(&self) -> MenteResult<()> {
585        let mut wal = self.wal.lock();
586        let mut pm = self.page_manager.lock();
587
588        wal.lock_exclusive()?;
589        wal.reload_lsn()?;
590
591        self.buffer_pool.flush_all(&mut pm)?;
592        pm.sync()?;
593
594        let lsn = wal.append(WalEntryType::Checkpoint, 0, &[])?;
595        wal.sync()?;
596        wal.truncate(lsn)?;
597        wal.unlock()?;
598
599        info!(lsn, "checkpoint complete");
600        Ok(())
601    }
602
603    /// Scan all pages and return (MemoryId, PageId) pairs for every valid memory node.
604    ///
605    /// Refreshes the page count from disk before scanning so pages written by
606    /// other processes are included. Used to rebuild the page map on startup.
607    ///
608    /// # Example
609    ///
610    /// ```no_run
611    /// # use mentedb_storage::StorageEngine;
612    /// # let engine = StorageEngine::open("/tmp/mentedb".as_ref())?;
613    /// let memories = engine.scan_all_memories();
614    /// for (memory_id, page_id) in &memories {
615    ///     println!("{memory_id} -> page {}", page_id.0);
616    /// }
617    /// # Ok::<(), mentedb_core::error::MenteError>(())
618    /// ```
619    pub fn scan_all_memories(&self) -> Vec<(mentedb_core::types::MemoryId, PageId)> {
620        let mut pm = self.page_manager.lock();
621        // Refresh from disk to see pages written by other processes
622        let _ = pm.reload_header();
623        let count = pm.page_count();
624        drop(pm);
625
626        let mut results = Vec::new();
627        for i in 1..count {
628            let page_id = PageId(i);
629            if let Ok(node) = self.load_memory(page_id) {
630                results.push((node.id, page_id));
631            }
632        }
633        results
634    }
635}
636
637#[cfg(test)]
638mod tests {
639    use super::*;
640    use mentedb_core::memory::MemoryType;
641    use mentedb_core::types::AgentId;
642
643    fn setup() -> (tempfile::TempDir, StorageEngine) {
644        let dir = tempfile::tempdir().unwrap();
645        let engine = StorageEngine::open(dir.path()).unwrap();
646        (dir, engine)
647    }
648
649    #[test]
650    fn test_allocate_write_read() {
651        let (_dir, engine) = setup();
652
653        let pid = engine.allocate_page().unwrap();
654        engine.write_page(pid, b"hello storage engine").unwrap();
655
656        let page = engine.read_page(pid).unwrap();
657        assert_eq!(&page.data[..20], b"hello storage engine");
658        engine.buffer_pool.unpin_page(pid, false).ok();
659    }
660
661    #[test]
662    fn test_store_and_load_memory() {
663        let (_dir, engine) = setup();
664
665        let node = MemoryNode::new(
666            AgentId::new(),
667            MemoryType::Episodic,
668            "The user prefers Rust over Go".to_string(),
669            vec![0.1, 0.2, 0.3, 0.4],
670        );
671
672        let page_id = engine.store_memory(&node).unwrap();
673        let loaded = engine.load_memory(page_id).unwrap();
674
675        assert_eq!(node.id, loaded.id);
676        assert_eq!(node.content, loaded.content);
677        assert_eq!(node.embedding, loaded.embedding);
678        assert_eq!(node.memory_type, loaded.memory_type);
679    }
680
681    #[test]
682    fn test_checkpoint() {
683        let (_dir, engine) = setup();
684
685        let node = MemoryNode::new(
686            AgentId::new(),
687            MemoryType::Semantic,
688            "checkpoint test".to_string(),
689            vec![1.0, 2.0],
690        );
691
692        let pid = engine.store_memory(&node).unwrap();
693        engine.checkpoint().unwrap();
694
695        let loaded = engine.load_memory(pid).unwrap();
696        assert_eq!(loaded.content, "checkpoint test");
697    }
698
699    #[test]
700    fn test_close_and_reopen() {
701        let dir = tempfile::tempdir().unwrap();
702        let pid;
703        {
704            let engine = StorageEngine::open(dir.path()).unwrap();
705            let node = MemoryNode::new(
706                AgentId::new(),
707                MemoryType::Procedural,
708                "persist across close".to_string(),
709                vec![0.5],
710            );
711            pid = engine.store_memory(&node).unwrap();
712            engine.close().unwrap();
713        }
714        {
715            let engine = StorageEngine::open(dir.path()).unwrap();
716            let loaded = engine.load_memory(pid).unwrap();
717            assert_eq!(loaded.content, "persist across close");
718        }
719    }
720
721    #[test]
722    fn test_crash_recovery() {
723        let dir = tempfile::tempdir().unwrap();
724        let mut ids = Vec::new();
725        let mut contents = Vec::new();
726        {
727            let engine = StorageEngine::open(dir.path()).unwrap();
728            for i in 0..3 {
729                let content = format!("crash-recovery-{i}");
730                let node = MemoryNode::new(
731                    AgentId::new(),
732                    MemoryType::Episodic,
733                    content.clone(),
734                    vec![i as f32],
735                );
736                let pid = engine.store_memory(&node).unwrap();
737                ids.push(pid);
738                contents.push(content);
739            }
740            // Simulate crash: sync the WAL but do NOT call close/checkpoint.
741            engine.wal.lock().sync().unwrap();
742        }
743        {
744            let engine = StorageEngine::open(dir.path()).unwrap();
745            for (pid, expected) in ids.iter().zip(contents.iter()) {
746                let loaded = engine.load_memory(*pid).unwrap();
747                assert_eq!(&loaded.content, expected);
748            }
749        }
750    }
751
752    #[test]
753    fn test_recovery_idempotent() {
754        let dir = tempfile::tempdir().unwrap();
755        let pid;
756        let content = "idempotent-check".to_string();
757        {
758            let engine = StorageEngine::open(dir.path()).unwrap();
759            let node = MemoryNode::new(
760                AgentId::new(),
761                MemoryType::Semantic,
762                content.clone(),
763                vec![1.0, 2.0],
764            );
765            pid = engine.store_memory(&node).unwrap();
766            engine.checkpoint().unwrap();
767            engine.close().unwrap();
768        }
769        {
770            let engine = StorageEngine::open(dir.path()).unwrap();
771            let loaded = engine.load_memory(pid).unwrap();
772            assert_eq!(loaded.content, content);
773        }
774    }
775
776    #[test]
777    fn test_partial_write_recovery() {
778        let dir = tempfile::tempdir().unwrap();
779        let mut ids = Vec::new();
780        let mut contents = Vec::new();
781        {
782            let engine = StorageEngine::open(dir.path()).unwrap();
783            for i in 0..3 {
784                let content = format!("checkpointed-{i}");
785                let node = MemoryNode::new(
786                    AgentId::new(),
787                    MemoryType::Semantic,
788                    content.clone(),
789                    vec![i as f32],
790                );
791                let pid = engine.store_memory(&node).unwrap();
792                ids.push(pid);
793                contents.push(content);
794            }
795            engine.checkpoint().unwrap();
796
797            for i in 3..5 {
798                let content = format!("unckeckpointed-{i}");
799                let node = MemoryNode::new(
800                    AgentId::new(),
801                    MemoryType::Episodic,
802                    content.clone(),
803                    vec![i as f32],
804                );
805                let pid = engine.store_memory(&node).unwrap();
806                ids.push(pid);
807                contents.push(content);
808            }
809            // Simulate crash — sync WAL but don't close.
810            engine.wal.lock().sync().unwrap();
811        }
812        {
813            let engine = StorageEngine::open(dir.path()).unwrap();
814            for (pid, expected) in ids.iter().zip(contents.iter()) {
815                let loaded = engine.load_memory(*pid).unwrap();
816                assert_eq!(&loaded.content, expected);
817            }
818        }
819    }
820
821    #[test]
822    fn test_delete_memory_durable() {
823        let dir = tempfile::tempdir().unwrap();
824        let pid;
825        {
826            let engine = StorageEngine::open(dir.path()).unwrap();
827            let node = MemoryNode::new(
828                AgentId::new(),
829                MemoryType::Semantic,
830                "to be deleted".to_string(),
831                vec![1.0],
832            );
833            pid = engine.store_memory(&node).unwrap();
834            engine.delete_memory(pid).unwrap();
835            assert!(engine.load_memory(pid).is_err());
836            assert!(engine.scan_all_memories().is_empty());
837            engine.close().unwrap();
838        }
839        {
840            let engine = StorageEngine::open(dir.path()).unwrap();
841            assert!(
842                engine.load_memory(pid).is_err(),
843                "deleted memory must not resurrect on reopen"
844            );
845            assert!(engine.scan_all_memories().is_empty());
846        }
847    }
848
849    #[test]
850    fn test_delete_survives_crash() {
851        let dir = tempfile::tempdir().unwrap();
852        let pid;
853        {
854            let engine = StorageEngine::open(dir.path()).unwrap();
855            let node = MemoryNode::new(
856                AgentId::new(),
857                MemoryType::Semantic,
858                "crash delete".to_string(),
859                vec![1.0],
860            );
861            pid = engine.store_memory(&node).unwrap();
862            engine.delete_memory(pid).unwrap();
863            // Simulate crash: no close, no checkpoint.
864        }
865        {
866            let engine = StorageEngine::open(dir.path()).unwrap();
867            assert!(
868                engine.load_memory(pid).is_err(),
869                "deletion must survive a crash via WAL replay"
870            );
871            assert!(engine.scan_all_memories().is_empty());
872        }
873    }
874
875    #[test]
876    fn test_deleted_page_reused() {
877        let (_dir, engine) = setup();
878
879        let a = MemoryNode::new(AgentId::new(), MemoryType::Semantic, "a".into(), vec![1.0]);
880        let pid_a = engine.store_memory(&a).unwrap();
881        engine.delete_memory(pid_a).unwrap();
882
883        let b = MemoryNode::new(AgentId::new(), MemoryType::Semantic, "b".into(), vec![2.0]);
884        let pid_b = engine.store_memory(&b).unwrap();
885        assert_eq!(pid_a, pid_b, "freed page should be reused");
886
887        let loaded = engine.load_memory(pid_b).unwrap();
888        assert_eq!(loaded.content, "b");
889    }
890
891    #[test]
892    fn test_delete_reuse_crash_recovery() {
893        let dir = tempfile::tempdir().unwrap();
894        let pid;
895        let b_id;
896        {
897            let engine = StorageEngine::open(dir.path()).unwrap();
898            let a = MemoryNode::new(AgentId::new(), MemoryType::Semantic, "a".into(), vec![1.0]);
899            pid = engine.store_memory(&a).unwrap();
900            engine.delete_memory(pid).unwrap();
901            let b = MemoryNode::new(AgentId::new(), MemoryType::Semantic, "b".into(), vec![2.0]);
902            let pid_b = engine.store_memory(&b).unwrap();
903            assert_eq!(pid, pid_b);
904            b_id = b.id;
905            // Simulate crash: the WAL now holds PageFree(p) then PageWrite(p).
906        }
907        {
908            let engine = StorageEngine::open(dir.path()).unwrap();
909            let loaded = engine.load_memory(pid).unwrap();
910            assert_eq!(loaded.content, "b", "later write must win over the free");
911            assert_eq!(loaded.id, b_id);
912            // The page must NOT be on the free list: a fresh allocation must
913            // not clobber b.
914            let c = MemoryNode::new(AgentId::new(), MemoryType::Semantic, "c".into(), vec![3.0]);
915            let pid_c = engine.store_memory(&c).unwrap();
916            assert_ne!(pid_c, pid, "recovered free list must exclude reused page");
917            assert_eq!(engine.load_memory(pid).unwrap().content, "b");
918        }
919    }
920
921    #[test]
922    fn test_update_memory_in_place() {
923        let dir = tempfile::tempdir().unwrap();
924        let pid;
925        let id;
926        {
927            let engine = StorageEngine::open(dir.path()).unwrap();
928            let mut node = MemoryNode::new(
929                AgentId::new(),
930                MemoryType::Semantic,
931                "original".to_string(),
932                vec![1.0],
933            );
934            pid = engine.store_memory(&node).unwrap();
935            id = node.id;
936
937            node.content = "updated".to_string();
938            engine.update_memory(pid, &node).unwrap();
939
940            let loaded = engine.load_memory(pid).unwrap();
941            assert_eq!(loaded.content, "updated");
942            // No orphan copy: exactly one entry in a full scan.
943            let scanned = engine.scan_all_memories();
944            assert_eq!(scanned.len(), 1);
945            engine.close().unwrap();
946        }
947        {
948            let engine = StorageEngine::open(dir.path()).unwrap();
949            let loaded = engine.load_memory(pid).unwrap();
950            assert_eq!(loaded.content, "updated");
951            assert_eq!(loaded.id, id);
952            assert_eq!(engine.scan_all_memories().len(), 1);
953        }
954    }
955
956    #[test]
957    fn test_concurrent_open_is_rejected() {
958        let dir = tempfile::tempdir().unwrap();
959
960        // The per operation flock serializes individual write transactions,
961        // but layers above the storage engine cache state across operations,
962        // so a second concurrent engine rolls the database back to its own
963        // open time snapshot. The process lock must refuse it.
964        let engine1 = StorageEngine::open(dir.path()).unwrap();
965        let second = StorageEngine::open(dir.path());
966        assert!(second.is_err(), "second concurrent open must fail");
967        let msg = second.err().unwrap().to_string();
968        assert!(msg.contains("locked"), "error names the lock: {msg}");
969
970        // Close releases the lock; the directory can be opened again.
971        engine1.close().unwrap();
972        let engine2 = StorageEngine::open(dir.path()).unwrap();
973        engine2.close().unwrap();
974    }
975
976    #[test]
977    fn test_concurrent_writes_from_threads() {
978        use std::sync::Arc;
979        let dir = tempfile::tempdir().unwrap();
980        let engine = Arc::new(StorageEngine::open(dir.path()).unwrap());
981
982        let handles: Vec<_> = (0..10)
983            .map(|i| {
984                let eng = Arc::clone(&engine);
985                std::thread::spawn(move || {
986                    let node = MemoryNode::new(
987                        AgentId::new(),
988                        MemoryType::Episodic,
989                        format!("thread-{i}"),
990                        vec![i as f32],
991                    );
992                    eng.store_memory(&node).unwrap()
993                })
994            })
995            .collect();
996
997        let pids: Vec<PageId> = handles.into_iter().map(|h| h.join().unwrap()).collect();
998
999        // All 10 writes succeeded and are readable
1000        for (i, pid) in pids.iter().enumerate() {
1001            let loaded = engine.load_memory(*pid).unwrap();
1002            assert_eq!(loaded.content, format!("thread-{i}"));
1003        }
1004    }
1005}