1use 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#[derive(Debug, Default)]
28pub struct TablePersistenceState {
29 pub columns: Vec<Column>,
31 pub flushed_rows: u64,
33 pub oversized: Vec<HashMap<u64, Vec<u8>>>,
37}
38
39#[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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}