Skip to main content

akar_storage/
persistence.rs

1//! Durable column mirror for in-memory tables (P45.4).
2//!
3//! Node/rel table rows are held in memory (`NodeGroup`s for node tables,
4//! edge/adjacency lists for rel tables). This module mirrors committed rows
5//! into persistent `Column` files (`col_{table_id}_{col_idx}` with a `.meta`
6//! sidecar) so that committed data survives process restarts.
7//!
8//! Lifecycle:
9//! - `sync_*` is called at commit/checkpoint. Newly inserted rows are appended
10//!   incrementally; when an UPDATE/DELETE touched a table, the mirror is
11//!   rewritten from scratch.
12//! - `load_*` is called at startup (`Database::new`) after tables are restored
13//!   from the persisted catalog, rebuilding the in-memory state from the mirror.
14//! - `remove` is called on DROP TABLE to delete the mirror files.
15
16use crate::buffer_manager::BufferManager;
17use crate::column::Column;
18use crate::table::{ColumnDefinition, NodeTable, RelTable, TableCatalog};
19use akar_common::enums::CompressionType;
20use akar_common::error::StorageError;
21use akar_common::types::{LogicalTypeID, Value};
22use std::collections::HashMap;
23use std::path::Path;
24use std::sync::{Arc, Mutex};
25
26/// Per-table durable mirror state.
27#[derive(Debug, Default)]
28pub struct TablePersistenceState {
29    /// Durable mirror columns, one per persisted table column.
30    pub columns: Vec<Column>,
31    /// Number of table rows already flushed into `columns`.
32    pub flushed_rows: u64,
33    /// Values larger than a single mirror page: `oversized[col][row]` holds the
34    /// serialised value bytes (the column keeps a `Null` placeholder so row
35    /// indices stay aligned). Persisted to the `col_{tid}.ovf` sidecar.
36    pub oversized: Vec<HashMap<u64, Vec<u8>>>,
37}
38
39/// Registry of durable column mirrors keyed by table id.
40#[derive(Debug, Default)]
41pub struct TablePersistence {
42    inner: Mutex<HashMap<u64, TablePersistenceState>>,
43}
44
45impl TablePersistence {
46    pub fn new() -> Self {
47        Self::default()
48    }
49
50    fn io_err(e: std::io::Error) -> StorageError {
51        StorageError::Page(format!("persistence I/O error: {e}"))
52    }
53
54    fn ovf_file_name(table_id: u64) -> String {
55        format!("col_{table_id}.ovf")
56    }
57
58    /// Append `value` to mirror column `col`. Values that are larger than a
59    /// single mirror page cannot be stored inline: they are recorded in
60    /// `oversized` (with a `Null` placeholder in the column so row indices stay
61    /// aligned) and persisted to the `.ovf` sidecar.
62    fn append_mirror_value(
63        col: &mut Column,
64        oversized: &mut HashMap<u64, Vec<u8>>,
65        row: u64,
66        value: &Value,
67    ) -> Result<(), StorageError> {
68        match col.append_value(value) {
69            Ok(()) => Ok(()),
70            Err(e) if e.kind() == std::io::ErrorKind::OutOfMemory => {
71                oversized.insert(row, Column::serialize_value(value));
72                col.append_value(&Value::Null).map_err(Self::io_err)
73            }
74            Err(e) => Err(Self::io_err(e)),
75        }
76    }
77
78    /// Persist the oversized-value map to the `col_{tid}.ovf` sidecar.
79    ///
80    /// Skipped entirely when no column has oversized values: most tables never
81    /// overflow, and rewriting (create + write + close) an empty sidecar on
82    /// every commit is wasted I/O on the hot commit path.
83    fn save_overflow(table_id: u64, oversized: &[HashMap<u64, Vec<u8>>], db_path: &Path) -> Result<(), StorageError> {
84        let path = db_path.join(Self::ovf_file_name(table_id));
85        if oversized.iter().all(|m| m.is_empty()) {
86            // No oversized values anywhere. The sidecar only exists if a past
87            // sync wrote one; drop it so recovery sees a consistent state.
88            if path.exists() {
89                let _ = std::fs::remove_file(&path);
90            }
91            return Ok(());
92        }
93        let mut buf: Vec<u8> = Vec::new();
94        buf.extend_from_slice(&(oversized.len() as u32).to_le_bytes());
95        for col in oversized {
96            let mut entries: Vec<(&u64, &Vec<u8>)> = col.iter().collect();
97            entries.sort_by_key(|(row, _)| **row);
98            buf.extend_from_slice(&(entries.len() as u32).to_le_bytes());
99            for (row, bytes) in entries {
100                buf.extend_from_slice(&row.to_le_bytes());
101                buf.extend_from_slice(&(bytes.len() as u32).to_le_bytes());
102                buf.extend_from_slice(bytes);
103            }
104        }
105        std::fs::write(&path, &buf).map_err(Self::io_err)
106    }
107
108    /// Load the oversized-value map from the `col_{tid}.ovf` sidecar.
109    fn load_overflow(table_id: u64, num_cols: usize, db_path: &Path) -> Vec<HashMap<u64, Vec<u8>>> {
110        let mut result: Vec<HashMap<u64, Vec<u8>>> = (0..num_cols).map(|_| HashMap::new()).collect();
111        let path = db_path.join(Self::ovf_file_name(table_id));
112        let buf = match std::fs::read(&path) {
113            Ok(b) => b,
114            Err(_) => return result,
115        };
116        let mut pos = 0usize;
117        if buf.len() < 4 {
118            return result;
119        }
120        let file_cols = u32::from_le_bytes(buf[pos..pos + 4].try_into().unwrap()) as usize;
121        pos += 4;
122        for ci in 0..file_cols.min(num_cols) {
123            if pos + 4 > buf.len() {
124                break;
125            }
126            let count = u32::from_le_bytes(buf[pos..pos + 4].try_into().unwrap()) as usize;
127            pos += 4;
128            for _ in 0..count {
129                if pos + 12 > buf.len() {
130                    break;
131                }
132                let row = u64::from_le_bytes(buf[pos..pos + 8].try_into().unwrap());
133                pos += 8;
134                let len = u32::from_le_bytes(buf[pos..pos + 4].try_into().unwrap()) as usize;
135                pos += 4;
136                if pos + len > buf.len() {
137                    break;
138                }
139                result[ci].insert(row, buf[pos..pos + len].to_vec());
140                pos += len;
141            }
142        }
143        result
144    }
145
146    fn build_column(
147        def: &ColumnDefinition,
148        table_id: u64,
149        col_idx: u32,
150        db_path: &Path,
151        bm: &Arc<Mutex<BufferManager>>,
152        page_size: usize,
153    ) -> Column {
154        Column::with_compression(
155            def.logical_type,
156            table_id,
157            col_idx,
158            db_path,
159            bm.clone(),
160            page_size,
161            def.compression,
162        )
163    }
164
165    /// Delete the mirror files for a dropped table.
166    pub fn remove(&self, table_id: u64, db_path: &Path, bm: &Arc<Mutex<BufferManager>>) {
167        if let Some(state) = self.inner.lock().unwrap().remove(&table_id) {
168            Self::drop_mirror_files(table_id, state.columns.len(), db_path, bm);
169        }
170    }
171
172    /// Remove the mirror files for columns `0..num_cols` of `table_id`.
173    ///
174    /// Drops the buffer-manager state (frames/mmap/registration) first so the
175    /// files can be deleted even while cached or memory-mapped, then removes
176    /// the data and `.meta` files from disk.
177    fn drop_mirror_files(table_id: u64, num_cols: usize, db_path: &Path, bm: &Arc<Mutex<BufferManager>>) {
178        {
179            let mut guard = bm.lock().unwrap();
180            for ci in 0..num_cols {
181                let fname = format!("col_{}_{}", table_id, ci);
182                guard.drop_file(&fname);
183            }
184        }
185        for ci in 0..num_cols {
186            let path = db_path.join(format!("col_{}_{}", table_id, ci));
187            let _ = std::fs::remove_file(&path);
188            let _ = std::fs::remove_file(path.with_extension("meta"));
189        }
190        let _ = std::fs::remove_file(db_path.join(Self::ovf_file_name(table_id)));
191    }
192
193    // ---------------------------------------------------------------------
194    // Node tables
195    // ---------------------------------------------------------------------
196
197    /// Persist `table`'s rows into its durable column mirror.
198    pub fn sync_node_table(
199        &self,
200        table: &mut NodeTable,
201        db_path: &Path,
202        bm: &Arc<Mutex<BufferManager>>,
203        page_size: usize,
204    ) -> Result<(), StorageError> {
205        let num_cols = table.columns.len();
206        let mut state = self.inner.lock().unwrap();
207        let entry = state.entry(table.table_id).or_default();
208
209        if !table.persistence_dirty && table.num_rows == entry.flushed_rows {
210            return Ok(());
211        }
212
213        if table.persistence_dirty || table.num_rows < entry.flushed_rows {
214            // Full rewrite: an UPDATE/DELETE touched the table, so rebuild the
215            // mirror from scratch (the in-memory state is the source of truth).
216            Self::drop_mirror_files(table.table_id, num_cols, db_path, bm);
217            let mut columns = Vec::with_capacity(num_cols);
218            let mut oversized: Vec<HashMap<u64, Vec<u8>>> = (0..num_cols).map(|_| HashMap::new()).collect();
219            for (ci, def) in table.columns.iter().enumerate() {
220                let mut col = Self::build_column(def, table.table_id, ci as u32, db_path, bm, page_size);
221                for row in 0..table.num_rows as usize {
222                    let value = table.get_value(row, ci).cloned().unwrap_or(Value::Null);
223                    Self::append_mirror_value(&mut col, &mut oversized[ci], row as u64, &value)?;
224                }
225                col.flush().map_err(Self::io_err)?;
226                col.save_metadata().map_err(Self::io_err)?;
227                columns.push(col);
228            }
229            entry.columns = columns;
230            entry.oversized = oversized;
231            entry.flushed_rows = table.num_rows;
232            table.persistence_dirty = false;
233        } else {
234            // Incremental append of newly inserted rows.
235            if entry.columns.is_empty() {
236                for (ci, def) in table.columns.iter().enumerate() {
237                    entry.columns.push(Self::build_column(
238                        def,
239                        table.table_id,
240                        ci as u32,
241                        db_path,
242                        bm,
243                        page_size,
244                    ));
245                }
246                entry.oversized = (0..num_cols).map(|_| HashMap::new()).collect();
247            }
248            for row in entry.flushed_rows as usize..table.num_rows as usize {
249                for (ci, col) in entry.columns.iter_mut().enumerate() {
250                    let value = table.get_value(row, ci).cloned().unwrap_or(Value::Null);
251                    Self::append_mirror_value(col, &mut entry.oversized[ci], row as u64, &value)?;
252                }
253            }
254            for col in entry.columns.iter_mut() {
255                col.flush().map_err(Self::io_err)?;
256                col.save_metadata().map_err(Self::io_err)?;
257            }
258            entry.flushed_rows = table.num_rows;
259        }
260        Self::save_overflow(table.table_id, &entry.oversized, db_path)?;
261        Ok(())
262    }
263
264    /// Load a node table from its durable column mirror.
265    ///
266    /// Returns `Ok(true)` when data was loaded, `Ok(false)` when no mirror
267    /// exists yet (fresh table).
268    pub fn load_node_table(
269        &self,
270        table: &mut NodeTable,
271        db_path: &Path,
272        bm: &Arc<Mutex<BufferManager>>,
273        page_size: usize,
274    ) -> Result<bool, StorageError> {
275        let num_cols = table.columns.len();
276        let mut columns = Vec::with_capacity(num_cols);
277        for (ci, def) in table.columns.iter().enumerate() {
278            let mut col = Self::build_column(def, table.table_id, ci as u32, db_path, bm, page_size);
279            if !col.load_metadata().map_err(Self::io_err)? {
280                return Ok(false);
281            }
282            columns.push(col);
283        }
284
285        let num_rows = columns[0].num_values as usize;
286        let oversized = Self::load_overflow(table.table_id, num_cols, db_path);
287        let mut rows = Vec::with_capacity(num_rows);
288        for row in 0..num_rows {
289            let mut values = Vec::with_capacity(num_cols);
290            for ci in 0..num_cols {
291                let value = if let Some(bytes) = oversized[ci].get(&(row as u64)) {
292                    Column::deserialize_value_bytes(bytes).unwrap_or(Value::Null)
293                } else {
294                    columns[ci].get_value(row as u64).unwrap_or(Value::Null)
295                };
296                values.push(value);
297            }
298            rows.push(values);
299        }
300
301        table.load_persisted_rows(rows)?;
302
303        let mut state = self.inner.lock().unwrap();
304        state.insert(
305            table.table_id,
306            TablePersistenceState {
307                columns,
308                flushed_rows: table.num_rows,
309                oversized,
310            },
311        );
312        Ok(num_rows > 0)
313    }
314
315    // ---------------------------------------------------------------------
316    // Rel tables
317    // ---------------------------------------------------------------------
318
319    /// Mirror column definition for the structural (src/dst) columns.
320    fn rel_structural_def(col_idx: usize) -> ColumnDefinition {
321        ColumnDefinition {
322            name: format!("__structural_{col_idx}"),
323            logical_type: LogicalTypeID::UInt64,
324            is_primary_key: false,
325            compression: CompressionType::Uncompressed,
326        }
327    }
328
329    /// Persist `table`'s edges + properties into its durable column mirror.
330    ///
331    /// Mirror layout: `[src: UInt64][dst: UInt64][prop_0..prop_n]`. Deleted
332    /// edges are stored as `(u64::MAX, u64::MAX)` tombstones so edge indices
333    /// stay stable across restarts.
334    pub fn sync_rel_table(
335        &self,
336        table: &mut RelTable,
337        db_path: &Path,
338        bm: &Arc<Mutex<BufferManager>>,
339        page_size: usize,
340    ) -> Result<(), StorageError> {
341        let num_prop_cols = table.columns.len();
342        let num_cols = num_prop_cols + 2;
343        let mut state = self.inner.lock().unwrap();
344        let entry = state.entry(table.table_id).or_default();
345
346        if !table.persistence_dirty && table.num_rows == entry.flushed_rows {
347            return Ok(());
348        }
349
350        let value_at = |table: &RelTable, ci: usize, e: usize| -> Value {
351            if ci == 0 {
352                Value::UInt64(table.edges[e].0)
353            } else if ci == 1 {
354                Value::UInt64(table.edges[e].1)
355            } else {
356                table.properties[ci - 2].get(e).cloned().unwrap_or(Value::Null)
357            }
358        };
359
360        if table.persistence_dirty || table.num_rows < entry.flushed_rows {
361            Self::drop_mirror_files(table.table_id, num_cols, db_path, bm);
362            let mut columns = Vec::with_capacity(num_cols);
363            let mut oversized: Vec<HashMap<u64, Vec<u8>>> = (0..num_cols).map(|_| HashMap::new()).collect();
364            for ci in 0..num_cols {
365                let def = if ci < 2 {
366                    Self::rel_structural_def(ci)
367                } else {
368                    table.columns[ci - 2].clone()
369                };
370                let mut col = Self::build_column(&def, table.table_id, ci as u32, db_path, bm, page_size);
371                for e in 0..table.num_rows as usize {
372                    Self::append_mirror_value(&mut col, &mut oversized[ci], e as u64, &value_at(table, ci, e))?;
373                }
374                col.flush().map_err(Self::io_err)?;
375                col.save_metadata().map_err(Self::io_err)?;
376                columns.push(col);
377            }
378            entry.columns = columns;
379            entry.oversized = oversized;
380            entry.flushed_rows = table.num_rows;
381            table.persistence_dirty = false;
382        } else {
383            if entry.columns.is_empty() {
384                for ci in 0..num_cols {
385                    let def = if ci < 2 {
386                        Self::rel_structural_def(ci)
387                    } else {
388                        table.columns[ci - 2].clone()
389                    };
390                    entry.columns.push(Self::build_column(
391                        &def,
392                        table.table_id,
393                        ci as u32,
394                        db_path,
395                        bm,
396                        page_size,
397                    ));
398                }
399                entry.oversized = (0..num_cols).map(|_| HashMap::new()).collect();
400            }
401            for e in entry.flushed_rows as usize..table.num_rows as usize {
402                for (ci, col) in entry.columns.iter_mut().enumerate() {
403                    Self::append_mirror_value(col, &mut entry.oversized[ci], e as u64, &value_at(table, ci, e))?;
404                }
405            }
406            for col in entry.columns.iter_mut() {
407                col.flush().map_err(Self::io_err)?;
408                col.save_metadata().map_err(Self::io_err)?;
409            }
410            entry.flushed_rows = table.num_rows;
411        }
412        Self::save_overflow(table.table_id, &entry.oversized, db_path)?;
413        Ok(())
414    }
415
416    /// Load a rel table from its durable column mirror.
417    pub fn load_rel_table(
418        &self,
419        table: &mut RelTable,
420        db_path: &Path,
421        bm: &Arc<Mutex<BufferManager>>,
422        page_size: usize,
423    ) -> Result<bool, StorageError> {
424        let num_prop_cols = table.columns.len();
425        let num_cols = num_prop_cols + 2;
426        let mut columns = Vec::with_capacity(num_cols);
427        for ci in 0..num_cols {
428            let def = if ci < 2 {
429                Self::rel_structural_def(ci)
430            } else {
431                table.columns[ci - 2].clone()
432            };
433            let mut col = Self::build_column(&def, table.table_id, ci as u32, db_path, bm, page_size);
434            if !col.load_metadata().map_err(Self::io_err)? {
435                return Ok(false);
436            }
437            columns.push(col);
438        }
439
440        let num_rows = columns[0].num_values as usize;
441        let oversized = Self::load_overflow(table.table_id, num_cols, db_path);
442        let mut edges = Vec::with_capacity(num_rows);
443        let mut properties = vec![Vec::with_capacity(num_rows); num_prop_cols];
444        let mut fwd_adj: HashMap<u64, Vec<(u64, usize)>> = HashMap::new();
445        let mut rev_adj: HashMap<u64, Vec<(u64, usize)>> = HashMap::new();
446        let structural = |ci: usize, e: u64| -> Value {
447            if let Some(bytes) = oversized[ci].get(&e) {
448                Column::deserialize_value_bytes(bytes).unwrap_or(Value::Null)
449            } else {
450                columns[ci].get_value(e).unwrap_or(Value::Null)
451            }
452        };
453        for e in 0..num_rows {
454            let src = match structural(0, e as u64) {
455                Value::UInt64(v) => v,
456                _ => u64::MAX,
457            };
458            let dst = match structural(1, e as u64) {
459                Value::UInt64(v) => v,
460                _ => u64::MAX,
461            };
462            edges.push((src, dst));
463            if src != u64::MAX {
464                fwd_adj.entry(src).or_default().push((dst, e));
465                rev_adj.entry(dst).or_default().push((src, e));
466            }
467            for ci in 0..num_prop_cols {
468                let prop = if let Some(bytes) = oversized[ci + 2].get(&(e as u64)) {
469                    Column::deserialize_value_bytes(bytes).unwrap_or(Value::Null)
470                } else {
471                    columns[ci + 2].get_value(e as u64).unwrap_or(Value::Null)
472                };
473                properties[ci].push(prop);
474            }
475        }
476
477        table.edges = edges;
478        table.fwd_adj = fwd_adj;
479        table.rev_adj = rev_adj;
480        table.properties = properties;
481        table.num_rows = num_rows as u64;
482        table.csr_index = None;
483
484        let mut state = self.inner.lock().unwrap();
485        state.insert(
486            table.table_id,
487            TablePersistenceState {
488                columns,
489                flushed_rows: table.num_rows,
490                oversized,
491            },
492        );
493        Ok(num_rows > 0)
494    }
495
496    // ---------------------------------------------------------------------
497    // All-tables helpers
498    // ---------------------------------------------------------------------
499
500    /// Sync all node + rel tables into their durable mirrors.
501    pub fn persist_all(
502        &self,
503        catalog: &Arc<TableCatalog>,
504        db_path: &Path,
505        bm: &Arc<Mutex<BufferManager>>,
506        page_size: usize,
507    ) -> Result<(), StorageError> {
508        let node_ids: Vec<u64> = catalog.all_node_tables().iter().map(|r| *r.key()).collect();
509        for tid in node_ids {
510            if let Some(mut table) = catalog.get_node_table_mut(tid) {
511                self.sync_node_table(&mut table, db_path, bm, page_size)?;
512            }
513        }
514        let rel_ids: Vec<u64> = catalog.all_rel_tables().iter().map(|r| *r.key()).collect();
515        for tid in rel_ids {
516            if let Some(mut table) = catalog.get_rel_table_mut(tid) {
517                self.sync_rel_table(&mut table, db_path, bm, page_size)?;
518            }
519        }
520        Ok(())
521    }
522
523    /// Load all persisted tables from their durable mirrors.
524    ///
525    /// Returns the number of tables that had persisted data.
526    pub fn load_all(
527        &self,
528        catalog: &Arc<TableCatalog>,
529        db_path: &Path,
530        bm: &Arc<Mutex<BufferManager>>,
531        page_size: usize,
532    ) -> Result<usize, StorageError> {
533        let mut loaded = 0usize;
534        let node_ids: Vec<u64> = catalog.all_node_tables().iter().map(|r| *r.key()).collect();
535        for tid in node_ids {
536            if let Some(mut table) = catalog.get_node_table_mut(tid) {
537                if self.load_node_table(&mut table, db_path, bm, page_size)? {
538                    loaded += 1;
539                }
540            }
541        }
542        let rel_ids: Vec<u64> = catalog.all_rel_tables().iter().map(|r| *r.key()).collect();
543        for tid in rel_ids {
544            if let Some(mut table) = catalog.get_rel_table_mut(tid) {
545                if self.load_rel_table(&mut table, db_path, bm, page_size)? {
546                    loaded += 1;
547                }
548            }
549        }
550        Ok(loaded)
551    }
552}
553
554#[cfg(test)]
555mod tests {
556    use super::*;
557    use crate::buffer_manager::{BufferManager, BufferManagerConfig};
558    use crate::table::{NodeTable, RelTable};
559    use akar_common::memory::MemoryManager;
560    use tempfile::TempDir;
561
562    fn test_dir() -> TempDir {
563        TempDir::new().expect("Failed to create temp dir")
564    }
565
566    fn buffer_manager(db_path: &Path) -> Arc<Mutex<BufferManager>> {
567        std::fs::create_dir_all(db_path).expect("Failed to create db dir");
568        let mm = Arc::new(MemoryManager::new(64 * 1024 * 1024));
569        Arc::new(Mutex::new(BufferManager::new(
570            db_path.to_path_buf(),
571            mm,
572            BufferManagerConfig::default(),
573        )))
574    }
575
576    fn node_defs() -> Vec<ColumnDefinition> {
577        vec![
578            ColumnDefinition {
579                name: "name".into(),
580                logical_type: LogicalTypeID::String,
581                is_primary_key: true,
582                compression: CompressionType::Uncompressed,
583            },
584            ColumnDefinition {
585                name: "age".into(),
586                logical_type: LogicalTypeID::Int64,
587                is_primary_key: false,
588                compression: CompressionType::Uncompressed,
589            },
590        ]
591    }
592
593    fn rel_defs() -> Vec<ColumnDefinition> {
594        vec![ColumnDefinition {
595            name: "since".into(),
596            logical_type: LogicalTypeID::Int64,
597            is_primary_key: false,
598            compression: CompressionType::Uncompressed,
599        }]
600    }
601
602    #[test]
603    fn test_node_table_mirror_roundtrip() {
604        let dir = test_dir();
605        let db_path = dir.path().join("db");
606        let bm = buffer_manager(&db_path);
607        let page_size = bm.lock().unwrap().page_size();
608        let persistence = TablePersistence::new();
609
610        let mut table = NodeTable::new(1, "Person".into(), node_defs());
611        table
612            .insert_row(vec![Value::String("alice".into()), Value::Int64(30)])
613            .unwrap();
614        table
615            .insert_row(vec![Value::String("bob".into()), Value::Int64(25)])
616            .unwrap();
617        table
618            .insert_row(vec![Value::String("carol".into()), Value::Int64(40)])
619            .unwrap();
620
621        persistence
622            .sync_node_table(&mut table, &db_path, &bm, page_size)
623            .unwrap();
624        assert!(!table.persistence_dirty, "sync should clear the dirty flag");
625
626        // Fresh table with the same schema must reload the mirrored rows.
627        let mut restored = NodeTable::new(1, "Person".into(), node_defs());
628        let loaded = persistence
629            .load_node_table(&mut restored, &db_path, &bm, page_size)
630            .unwrap();
631        assert!(loaded, "mirror should exist and be loadable");
632        assert_eq!(restored.num_rows, 3);
633        assert_eq!(restored.get_value(0, 0), Some(&Value::String("alice".into())));
634        assert_eq!(restored.get_value(1, 1), Some(&Value::Int64(25)));
635        assert_eq!(restored.get_value(2, 0), Some(&Value::String("carol".into())));
636    }
637
638    #[test]
639    fn test_node_table_mirror_incremental_append() {
640        let dir = test_dir();
641        let db_path = dir.path().join("db");
642        let bm = buffer_manager(&db_path);
643        let page_size = bm.lock().unwrap().page_size();
644        let persistence = TablePersistence::new();
645
646        let mut table = NodeTable::new(1, "Person".into(), node_defs());
647        table
648            .insert_row(vec![Value::String("alice".into()), Value::Int64(30)])
649            .unwrap();
650        persistence
651            .sync_node_table(&mut table, &db_path, &bm, page_size)
652            .unwrap();
653
654        // Insert more rows — the mirror must append incrementally.
655        table
656            .insert_row(vec![Value::String("bob".into()), Value::Int64(25)])
657            .unwrap();
658        persistence
659            .sync_node_table(&mut table, &db_path, &bm, page_size)
660            .unwrap();
661
662        let mut restored = NodeTable::new(1, "Person".into(), node_defs());
663        persistence
664            .load_node_table(&mut restored, &db_path, &bm, page_size)
665            .unwrap();
666        assert_eq!(restored.num_rows, 2);
667        assert_eq!(restored.get_value(1, 0), Some(&Value::String("bob".into())));
668    }
669
670    #[test]
671    fn test_node_table_mirror_update_triggers_rewrite() {
672        let dir = test_dir();
673        let db_path = dir.path().join("db");
674        let bm = buffer_manager(&db_path);
675        let page_size = bm.lock().unwrap().page_size();
676        let persistence = TablePersistence::new();
677
678        let mut table = NodeTable::new(1, "Person".into(), node_defs());
679        table
680            .insert_row(vec![Value::String("alice".into()), Value::Int64(30)])
681            .unwrap();
682        persistence
683            .sync_node_table(&mut table, &db_path, &bm, page_size)
684            .unwrap();
685
686        // UPDATE marks the table dirty → the mirror is rewritten from scratch.
687        table.update_cell(0, 1, Value::Int64(31)).unwrap();
688        assert!(table.persistence_dirty);
689        persistence
690            .sync_node_table(&mut table, &db_path, &bm, page_size)
691            .unwrap();
692        assert!(!table.persistence_dirty);
693
694        let mut restored = NodeTable::new(1, "Person".into(), node_defs());
695        persistence
696            .load_node_table(&mut restored, &db_path, &bm, page_size)
697            .unwrap();
698        assert_eq!(restored.num_rows, 1);
699        assert_eq!(restored.get_value(0, 1), Some(&Value::Int64(31)));
700    }
701
702    #[test]
703    fn test_rel_table_mirror_roundtrip() {
704        let dir = test_dir();
705        let db_path = dir.path().join("db");
706        let bm = buffer_manager(&db_path);
707        let page_size = bm.lock().unwrap().page_size();
708        let persistence = TablePersistence::new();
709
710        let mut table = RelTable::new(2, "LivesIn".into(), 1, 1, rel_defs());
711        table.insert_rel(0, 1, vec![Value::Int64(2010)]).unwrap();
712        table.insert_rel(0, 2, vec![Value::Int64(2015)]).unwrap();
713        table.insert_rel(3, 1, vec![Value::Int64(2020)]).unwrap();
714
715        persistence
716            .sync_rel_table(&mut table, &db_path, &bm, page_size)
717            .unwrap();
718        assert!(!table.persistence_dirty);
719
720        let mut restored = RelTable::new(2, "LivesIn".into(), 1, 1, rel_defs());
721        let loaded = persistence
722            .load_rel_table(&mut restored, &db_path, &bm, page_size)
723            .unwrap();
724        assert!(loaded, "rel mirror should exist and be loadable");
725        assert_eq!(restored.num_rows, 3);
726        assert_eq!(restored.edges, vec![(0, 1), (0, 2), (3, 1)]);
727        assert_eq!(restored.fwd_adj.get(&0).cloned(), Some(vec![(1, 0), (2, 1)]));
728        assert_eq!(restored.rev_adj.get(&1).cloned(), Some(vec![(0, 0), (3, 2)]));
729        assert_eq!(
730            restored.properties[0],
731            vec![Value::Int64(2010), Value::Int64(2015), Value::Int64(2020)]
732        );
733    }
734
735    #[test]
736    fn test_node_table_mirror_oversized_value() {
737        let dir = test_dir();
738        let db_path = dir.path().join("db");
739        let bm = buffer_manager(&db_path);
740        let page_size = bm.lock().unwrap().page_size();
741        let persistence = TablePersistence::new();
742
743        let long = "A".repeat(page_size * 4);
744        let mut table = NodeTable::new(1, "Person".into(), node_defs());
745        table
746            .insert_row(vec![Value::String("alice".into()), Value::Int64(30)])
747            .unwrap();
748        table
749            .insert_row(vec![Value::String(long.clone()), Value::Int64(25)])
750            .unwrap();
751        table
752            .insert_row(vec![Value::String("carol".into()), Value::Int64(40)])
753            .unwrap();
754
755        persistence
756            .sync_node_table(&mut table, &db_path, &bm, page_size)
757            .unwrap();
758        assert!(!table.persistence_dirty);
759
760        // The oversized string is stored in the overflow sidecar.
761        assert!(db_path.join("col_1.ovf").exists(), "overflow sidecar should be written");
762
763        let mut restored = NodeTable::new(1, "Person".into(), node_defs());
764        persistence
765            .load_node_table(&mut restored, &db_path, &bm, page_size)
766            .unwrap();
767        assert_eq!(restored.num_rows, 3);
768        assert_eq!(restored.get_value(0, 0), Some(&Value::String("alice".into())));
769        assert_eq!(restored.get_value(1, 0), Some(&Value::String(long.clone())));
770        assert_eq!(restored.get_value(1, 1), Some(&Value::Int64(25)));
771        assert_eq!(restored.get_value(2, 0), Some(&Value::String("carol".into())));
772
773        // A subsequent incremental append must keep the oversized row intact.
774        table
775            .insert_row(vec![Value::String("dave".into()), Value::Int64(50)])
776            .unwrap();
777        persistence
778            .sync_node_table(&mut table, &db_path, &bm, page_size)
779            .unwrap();
780        let mut restored2 = NodeTable::new(1, "Person".into(), node_defs());
781        persistence
782            .load_node_table(&mut restored2, &db_path, &bm, page_size)
783            .unwrap();
784        assert_eq!(restored2.num_rows, 4);
785        assert_eq!(restored2.get_value(1, 0), Some(&Value::String(long.clone())));
786        assert_eq!(restored2.get_value(3, 0), Some(&Value::String("dave".into())));
787    }
788
789    #[test]
790    fn test_remove_deletes_mirror_files() {
791        let dir = test_dir();
792        let db_path = dir.path().join("db");
793        let bm = buffer_manager(&db_path);
794        let page_size = bm.lock().unwrap().page_size();
795        let persistence = TablePersistence::new();
796
797        let mut table = NodeTable::new(1, "Person".into(), node_defs());
798        table
799            .insert_row(vec![Value::String("alice".into()), Value::Int64(30)])
800            .unwrap();
801        persistence
802            .sync_node_table(&mut table, &db_path, &bm, page_size)
803            .unwrap();
804
805        let col_file = db_path.join("col_1_0");
806        assert!(col_file.exists(), "column data file should exist after sync");
807
808        persistence.remove(1, &db_path, &bm);
809        assert!(!col_file.exists(), "column data file should be removed on drop");
810        assert!(!db_path.join("col_1_0.meta").exists());
811    }
812}