1use crate::art_index::ArtPrimaryKeyIndex;
4use crate::art_key::ArtKey;
5use crate::column::Column;
6use crate::column_chunk::ColumnChunk;
7use crate::csr::CsrIndex;
8use crate::index::HashIndex;
9use crate::node_group::NodeGroup;
10use crate::spiller::Spiller;
11use crate::vector_index::VectorIndexTable;
12use akar_common::error::StorageError;
13use akar_common::types::{LogicalTypeID, Value, pk_value_to_string};
14use akar_vector::hnsw::DistanceMetric;
15use dashmap::DashMap;
16use std::collections::HashMap;
17use std::sync::Arc;
18
19#[derive(Debug, Clone)]
21pub struct ColumnDefinition {
22 pub name: String,
23 pub logical_type: LogicalTypeID,
24 pub is_primary_key: bool,
25 pub compression: akar_common::enums::CompressionType,
26}
27
28#[derive(Debug, Clone)]
39pub struct NodeTable {
40 pub table_id: u64,
41 pub name: String,
42 pub columns: Vec<ColumnDefinition>,
43 pub primary_key_column: usize,
44 pub num_rows: u64,
45 pub node_groups: Vec<NodeGroup>,
48 pub hash_index: HashIndex<String>,
51 pub art_index: Option<ArtPrimaryKeyIndex>,
54 pub persistence_dirty: bool,
57 spiller: Option<Arc<Spiller>>,
62}
63
64pub const NO_PRIMARY_KEY: usize = usize::MAX;
70
71impl NodeTable {
72 pub fn new(table_id: u64, name: String, columns: Vec<ColumnDefinition>) -> Self {
73 let primary_key_column = columns.iter().position(|c| c.is_primary_key).unwrap_or(NO_PRIMARY_KEY);
78 Self {
79 table_id,
80 name,
81 columns,
82 primary_key_column,
83 num_rows: 0,
84 node_groups: Vec::new(),
85 hash_index: HashIndex::new(),
86 art_index: None,
87 persistence_dirty: false,
88 spiller: None,
89 }
90 }
91
92 pub fn set_spiller(&mut self, spiller: Option<Arc<Spiller>>) {
95 self.spiller = spiller;
96 }
97
98 pub fn add_column(&mut self, column: ColumnDefinition) {
104 if self.columns.iter().any(|c| c.name.eq_ignore_ascii_case(&column.name)) {
105 return;
106 }
107 self.columns.push(column);
108 for group in &mut self.node_groups {
109 let mut chunk = ColumnChunk::new();
110 for _ in 0..group.num_nodes {
111 chunk.append(Value::Null);
112 }
113 group.columns.push(chunk);
114 }
115 self.persistence_dirty = true;
116 }
117
118 pub fn insert_row(&mut self, values: Vec<Value>) -> Result<u64, StorageError> {
133 self.insert_row_with_txn(values, None)
134 }
135
136 pub fn insert_row_with_txn(&mut self, mut values: Vec<Value>, txn_id: Option<u64>) -> Result<u64, StorageError> {
138 if values.len() != self.columns.len() {
139 return Err(StorageError::Page(format!(
140 "Column count mismatch: expected {} values, got {}",
141 self.columns.len(),
142 values.len()
143 )));
144 }
145
146 if self.primary_key_column < self.columns.len() {
148 let pk_value = &values[self.primary_key_column];
149 if matches!(pk_value, Value::Null) {
150 return Err(StorageError::Page(format!(
151 "NULL value not allowed for primary key column '{}' in table '{}'",
152 self.columns[self.primary_key_column].name, self.name
153 )));
154 }
155 }
156
157 coerce_values_to_columns(&mut values, &self.columns)?;
162
163 if self.primary_key_column < self.columns.len() {
165 let pk_value = &values[self.primary_key_column];
166 let pk_key = pk_value_to_string(pk_value);
167 if self.hash_index.lookup(&pk_key).is_some() {
168 return Err(StorageError::Index(format!(
169 "Duplicate primary key value: '{pk_key}' in table '{}'",
170 self.name
171 )));
172 }
173 }
174
175 let num_cols = self.columns.len();
177 if self.node_groups.is_empty() || self.node_groups.last().unwrap().is_full() {
178 let start_offset = self.num_rows;
179 let mut new_group = NodeGroup::new(num_cols, start_offset);
180 if txn_id.is_some() {
182 new_group.enable_version_info();
183 }
184 if let Some(ref spiller) = self.spiller {
185 new_group.set_spiller(spiller.clone());
186 }
187 self.node_groups.push(new_group);
188 }
189
190 let current = self.node_groups.last_mut().unwrap();
191 if txn_id.is_some() {
193 current.enable_version_info();
194 }
195 current.append_row_with_txn(values.clone(), txn_id)?;
196 current.restore_spilled()?;
199 self.num_rows += 1;
200
201 self.persistence_dirty = true;
205
206 if self.primary_key_column < self.columns.len() {
208 let pk_value = &values[self.primary_key_column];
209 let pk_key = pk_value_to_string(pk_value);
210 self.hash_index.insert(pk_key, self.num_rows - 1);
211
212 if let Some(ref mut art_idx) = self.art_index
214 && let Some(art_key) = ArtKey::from_value(pk_value)
215 {
216 art_idx.insert(&art_key, self.num_rows - 1);
217 }
218 }
219
220 Ok(self.num_rows - 1)
221 }
222
223 pub fn insert_rows_batch(&mut self, rows: &[Vec<Value>]) -> Result<u64, StorageError> {
227 self.insert_rows_batch_with_txn(rows, None)
228 }
229
230 pub fn insert_rows_batch_with_txn(
232 &mut self,
233 rows: &[Vec<Value>],
234 txn_id: Option<u64>,
235 ) -> Result<u64, StorageError> {
236 if rows.is_empty() {
237 return Ok(0);
238 }
239 let num_cols = self.columns.len();
240
241 for (i, row) in rows.iter().enumerate() {
243 if row.len() != num_cols {
244 return Err(StorageError::Page(format!(
245 "Row {} column count mismatch: expected {} values, got {}",
246 i,
247 num_cols,
248 row.len()
249 )));
250 }
251 if self.primary_key_column < num_cols {
253 let pk_value = &row[self.primary_key_column];
254 if matches!(pk_value, Value::Null) {
255 return Err(StorageError::Page(format!(
256 "NULL value not allowed for primary key column '{}' in table '{}'",
257 self.columns[self.primary_key_column].name, self.name
258 )));
259 }
260 }
261 if self.primary_key_column < num_cols {
263 let pk_key = pk_value_to_string(&row[self.primary_key_column]);
264 if self.hash_index.lookup(&pk_key).is_some() {
265 return Err(StorageError::Index(format!(
266 "Duplicate primary key value: '{pk_key}' in table '{}'",
267 self.name
268 )));
269 }
270 }
271 }
272
273 let start_offset = self.num_rows;
274 let total_new = rows.len();
275
276 let mut coerced_rows = rows.to_vec();
280 for row in &mut coerced_rows {
281 coerce_values_to_columns(row, &self.columns)?;
282 }
283 let rows: &[Vec<Value>] = &coerced_rows;
284
285 if self.node_groups.is_empty() || self.node_groups.last().unwrap().is_full() {
287 let off = if self.node_groups.is_empty() {
288 start_offset
289 } else {
290 self.num_rows
291 };
292 let mut new_group = NodeGroup::new(num_cols, off);
293 if txn_id.is_some() {
294 new_group.enable_version_info();
295 }
296 if let Some(ref spiller) = self.spiller {
297 new_group.set_spiller(spiller.clone());
298 }
299 self.node_groups.push(new_group);
300 }
301
302 let mut inserted = 0usize;
304 while inserted < total_new {
305 let current = self.node_groups.last_mut().unwrap();
306 if txn_id.is_some() {
307 current.enable_version_info();
308 }
309 let rem = current.remaining();
310 let take = (total_new - inserted).min(rem);
311 for row in &rows[inserted..inserted + take] {
312 current.append_row_with_txn(row.clone(), txn_id)?;
313 }
314 self.num_rows += take as u64;
315 inserted += take;
316 if inserted < total_new {
317 let off = self.num_rows;
318 let mut new_group = NodeGroup::new(num_cols, off);
319 if txn_id.is_some() {
320 new_group.enable_version_info();
321 }
322 if let Some(ref spiller) = self.spiller {
323 new_group.set_spiller(spiller.clone());
324 }
325 self.node_groups.push(new_group);
326 }
327 }
328
329 for group in &mut self.node_groups {
332 group.restore_spilled()?;
333 }
334
335 self.persistence_dirty = true;
339
340 for (i, row) in rows.iter().enumerate() {
342 if self.primary_key_column < num_cols {
343 let pk_key = pk_value_to_string(&row[self.primary_key_column]);
344 self.hash_index.insert(pk_key, start_offset + i as u64);
345 if let Some(ref mut art_idx) = self.art_index
346 && let Some(art_key) = ArtKey::from_value(&row[self.primary_key_column])
347 {
348 art_idx.insert(&art_key, start_offset + i as u64);
349 }
350 }
351 }
352
353 Ok(rows.len() as u64)
354 }
355
356 pub fn lookup_by_pk(&self, pk_value: &Value) -> Option<u64> {
361 let pk_key = pk_value_to_string(pk_value);
362 self.hash_index.lookup(&pk_key)
363 }
364
365 pub fn lookup_by_pk_batch(&self, pk_values: &[Value]) -> Vec<Option<u64>> {
372 pk_values
373 .iter()
374 .map(|pk_value| {
375 let pk_key = pk_value_to_string(pk_value);
376 self.hash_index.lookup(&pk_key)
377 })
378 .collect()
379 }
380
381 pub fn lookup_by_pk_range(
389 &self,
390 lower: Option<&Value>,
391 lower_inclusive: bool,
392 upper: Option<&Value>,
393 upper_inclusive: bool,
394 max_results: u64,
395 ) -> Vec<u64> {
396 match &self.art_index {
397 Some(idx) => {
398 let lower_key = lower.and_then(ArtKey::from_value);
399 let upper_key = upper.and_then(ArtKey::from_value);
400 idx.range_scan(
401 lower_key.as_ref(),
402 lower_inclusive,
403 upper_key.as_ref(),
404 upper_inclusive,
405 max_results,
406 )
407 }
408 None => Vec::new(),
409 }
410 }
411
412 pub fn scan_column(
421 &self,
422 col_idx: usize,
423 start: u64,
424 count: u64,
425 snapshot_ts: Option<u64>,
426 commit_history: &HashMap<u64, u64>,
427 ) -> Vec<Value> {
428 if col_idx >= self.columns.len() || start >= self.num_rows {
429 return Vec::new();
430 }
431 let end = (start + count).min(self.num_rows);
432 let mut result = Vec::with_capacity((end - start) as usize);
433
434 let group_start = self.find_group(start);
436 let mut remaining = end - start;
437
438 for g_idx in group_start..self.node_groups.len() {
439 if remaining == 0 {
440 break;
441 }
442 let group = &self.node_groups[g_idx];
443 let local_start = if g_idx == group_start {
444 (start - group.start_offset) as usize
445 } else {
446 0
447 };
448 let available = (group.num_nodes as usize).saturating_sub(local_start);
449 let take = available.min(remaining as usize);
450
451 for row in local_start..local_start + take {
452 let val = group.get_value_with_snapshot(row, col_idx, snapshot_ts, commit_history);
453 match val {
454 Some(v) => result.push(v.clone()),
455 None => result.push(Value::Null),
456 }
457 }
458 remaining -= take as u64;
459 }
460
461 result
462 }
463
464 pub fn load_persisted_rows(&mut self, rows: Vec<Vec<Value>>) -> Result<(), StorageError> {
472 let num_cols = self.columns.len();
473 self.node_groups.clear();
474 self.hash_index.clear();
475
476 let mut offset = 0u64;
477 let mut group = NodeGroup::new(num_cols, offset);
478 for row in &rows {
479 if row.len() != num_cols {
480 return Err(StorageError::Page(format!(
481 "load_persisted_rows: expected {num_cols} values, got {}",
482 row.len()
483 )));
484 }
485 group.append_row_with_txn(row.clone(), None)?;
486 offset += 1;
487 if group.is_full() {
488 self.node_groups.push(group);
489 group = NodeGroup::new(num_cols, offset);
490 }
491 }
492 if group.num_nodes > 0 {
493 self.node_groups.push(group);
494 }
495 self.num_rows = offset;
496
497 if self.primary_key_column < num_cols {
499 for (row_idx, values) in rows.iter().enumerate() {
500 let pk = &values[self.primary_key_column];
501 if matches!(pk, Value::Null) {
502 continue;
503 }
504 let key = pk_value_to_string(pk);
505 self.hash_index.insert(key, row_idx as u64);
506 }
507 }
508
509 if self.art_index.is_some() && self.primary_key_column < num_cols {
511 if let Some(art) = &mut self.art_index {
512 art.clear();
513 for (row_idx, values) in rows.iter().enumerate() {
514 let pk = &values[self.primary_key_column];
515 if matches!(pk, Value::Null) {
516 continue;
517 }
518 if let Some(key) = ArtKey::from_value(pk) {
519 art.insert(&key, row_idx as u64);
520 }
521 }
522 }
523 }
524 Ok(())
525 }
526
527 pub fn update_cell(&mut self, row_idx: u64, col_idx: usize, value: Value) -> Result<(), StorageError> {
529 if col_idx >= self.columns.len() {
530 return Err(StorageError::Page(format!("Column index {col_idx} out of range")));
531 }
532 if row_idx >= self.num_rows {
533 return Err(StorageError::Page(format!(
534 "Row index {row_idx} out of range (num_rows={})",
535 self.num_rows
536 )));
537 }
538 self.persistence_dirty = true;
539
540 let mut offset = 0u64;
541 for group in &mut self.node_groups {
542 if row_idx < offset + group.num_nodes {
543 let local_row = (row_idx - offset) as usize;
544 if let Some(col_chunk) = group.columns.get_mut(col_idx) {
545 col_chunk.set_value(local_row, value)?;
546 }
547 return Ok(());
548 }
549 offset += group.num_nodes;
550 }
551 Err(StorageError::Page(format!(
552 "Row index {row_idx} not found in any node group"
553 )))
554 }
555
556 pub fn delete_row(&mut self, row_idx: u64) -> Result<(), StorageError> {
559 self.delete_row_with_txn(row_idx, None)
560 }
561
562 pub fn delete_row_with_txn(&mut self, row_idx: u64, txn_id: Option<u64>) -> Result<(), StorageError> {
567 if row_idx >= self.num_rows {
568 return Err(StorageError::Page(format!(
569 "Row index {row_idx} out of range (num_rows={})",
570 self.num_rows
571 )));
572 }
573 self.persistence_dirty = true;
574
575 let mut offset = 0u64;
577 for group in &mut self.node_groups {
578 if row_idx < offset + group.num_nodes {
579 let local_row = (row_idx - offset) as usize;
580 if let Some(txn) = txn_id {
582 if let Some(ref vi) = group.version_info {
583 vi.delete(txn, local_row as u32);
584 }
585 }
586 let pk_value = (self.primary_key_column < self.columns.len())
589 .then(|| group.columns.get(self.primary_key_column))
590 .flatten()
591 .and_then(|chunk| chunk.get(local_row))
592 .cloned();
593 for col_chunk in &mut group.columns {
595 let _ = col_chunk.set_value(local_row, Value::Null);
596 }
597 if let Some(pk) = pk_value
601 && !matches!(pk, Value::Null)
602 {
603 let pk_key = pk_value_to_string(&pk);
604 self.hash_index.delete(&pk_key);
605 if let Some(ref mut art_idx) = self.art_index
606 && let Some(art_key) = ArtKey::from_value(&pk)
607 {
608 art_idx.delete(&art_key, row_idx);
609 }
610 }
611 return Ok(());
612 }
613 offset += group.num_nodes;
614 }
615 Err(StorageError::Page(format!(
616 "Row index {row_idx} not found in any node group"
617 )))
618 }
619
620 pub fn row_undo_bytes(&self, row_idx: u64) -> Vec<u8> {
624 let mut out = Vec::new();
625 for col in 0..self.columns.len() {
626 let val = self.get_value(row_idx as usize, col).cloned().unwrap_or(Value::Null);
627 out.extend_from_slice(&Column::serialize_value(&val));
628 }
629 out
630 }
631
632 pub fn cell_undo_bytes(&self, row_idx: u64, col_idx: usize) -> Vec<u8> {
635 let val = self
636 .get_value(row_idx as usize, col_idx)
637 .cloned()
638 .unwrap_or(Value::Null);
639 Column::serialize_value(&val)
640 }
641
642 pub fn get_value(&self, row: usize, col: usize) -> Option<&Value> {
645 self.get_value_with_snapshot(row, col, None, &HashMap::new())
646 }
647
648 pub fn get_value_with_snapshot(
653 &self,
654 row: usize,
655 col: usize,
656 snapshot_ts: Option<u64>,
657 commit_history: &HashMap<u64, u64>,
658 ) -> Option<&Value> {
659 if col >= self.columns.len() || row as u64 >= self.num_rows {
660 return None;
661 }
662 let group_idx = self.find_group(row as u64);
663 let group = self.node_groups.get(group_idx)?;
664 let local_row = row as u64 - group.start_offset;
665 group.get_value_with_snapshot(local_row as usize, col, snapshot_ts, commit_history)
666 }
667
668 pub fn to_column_major_data(&self) -> Vec<Vec<Value>> {
672 self.to_column_major_data_with_predicate(None)
673 }
674
675 pub fn to_column_major_data_with_predicate(&self, predicate: Option<(usize, &str, &Value)>) -> Vec<Vec<Value>> {
678 let num_cols = self.columns.len();
679 let mut result = vec![Vec::new(); num_cols]; for group in &self.node_groups {
682 if let Some((col_idx, op, val)) = predicate
683 && let Some(col_chunk) = group.columns.get(col_idx)
684 {
685 use crate::predicate::{ZoneMapCheckResult, check_zone_map};
686 if check_zone_map(&col_chunk.stats, op, val) == ZoneMapCheckResult::SkipScan {
687 continue; }
689 }
690
691 for row in 0..group.num_nodes as usize {
692 for (col, res_col) in result.iter_mut().enumerate().take(num_cols) {
693 match group.get_value(row, col) {
694 Some(v) => res_col.push(v.clone()),
695 None => res_col.push(Value::Null),
696 }
697 }
698 }
699 }
700
701 result
702 }
703
704 pub fn to_column_major_data_with_snapshot(
710 &self,
711 snapshot_ts: Option<u64>,
712 commit_history: &HashMap<u64, u64>,
713 ) -> Vec<Vec<Value>> {
714 let num_cols = self.columns.len();
715 let mut result = vec![Vec::new(); num_cols];
716
717 for group in &self.node_groups {
718 for row in 0..group.num_nodes as usize {
719 for (col, res_col) in result.iter_mut().enumerate().take(num_cols) {
720 match group.get_value_owned_with_snapshot(row, col, snapshot_ts, commit_history) {
721 Some(v) => res_col.push(v),
722 None => res_col.push(Value::Null),
723 }
724 }
725 }
726 }
727
728 result
729 }
730
731 pub fn to_column_major_data_with_snapshot_and_predicate(
734 &self,
735 predicate: Option<(usize, &str, &Value)>,
736 snapshot_ts: Option<u64>,
737 commit_history: &HashMap<u64, u64>,
738 ) -> Vec<Vec<Value>> {
739 let num_cols = self.columns.len();
740 let mut result = vec![Vec::new(); num_cols];
741
742 for group in &self.node_groups {
743 if let Some((col_idx, op, val)) = predicate
744 && let Some(col_chunk) = group.columns.get(col_idx)
745 {
746 use crate::predicate::{ZoneMapCheckResult, check_zone_map};
747 if check_zone_map(&col_chunk.stats, op, val) == ZoneMapCheckResult::SkipScan {
748 continue;
749 }
750 }
751
752 for row in 0..group.num_nodes as usize {
753 for (col, res_col) in result.iter_mut().enumerate().take(num_cols) {
754 match group.get_value_owned_with_snapshot(row, col, snapshot_ts, commit_history) {
755 Some(v) => res_col.push(v),
756 None => res_col.push(Value::Null),
757 }
758 }
759 }
760 }
761
762 result
763 }
764
765 pub fn to_column_major_data_with_predicate_and_ids(
770 &self,
771 predicate: Option<(usize, &str, &Value)>,
772 ) -> (Vec<Vec<Value>>, Vec<u64>) {
773 let num_cols = self.columns.len();
774 let mut result = vec![Vec::new(); num_cols];
775 let mut ids = Vec::new();
776
777 for group in &self.node_groups {
778 if let Some((col_idx, op, val)) = predicate
779 && let Some(col_chunk) = group.columns.get(col_idx)
780 {
781 use crate::predicate::{ZoneMapCheckResult, check_zone_map};
782 if check_zone_map(&col_chunk.stats, op, val) == ZoneMapCheckResult::SkipScan {
783 continue;
784 }
785 }
786
787 for row in 0..group.num_nodes as usize {
788 ids.push(group.start_offset + row as u64);
789 for (col, res_col) in result.iter_mut().enumerate().take(num_cols) {
790 match group.get_value(row, col) {
791 Some(v) => res_col.push(v.clone()),
792 None => res_col.push(Value::Null),
793 }
794 }
795 }
796 }
797
798 (result, ids)
799 }
800
801 pub fn to_column_major_data_with_snapshot_and_predicate_and_ids(
805 &self,
806 predicate: Option<(usize, &str, &Value)>,
807 snapshot_ts: Option<u64>,
808 commit_history: &HashMap<u64, u64>,
809 ) -> (Vec<Vec<Value>>, Vec<u64>) {
810 let num_cols = self.columns.len();
811 let mut result = vec![Vec::new(); num_cols];
812 let mut ids = Vec::new();
813
814 for group in &self.node_groups {
815 if let Some((col_idx, op, val)) = predicate
816 && let Some(col_chunk) = group.columns.get(col_idx)
817 {
818 use crate::predicate::{ZoneMapCheckResult, check_zone_map};
819 if check_zone_map(&col_chunk.stats, op, val) == ZoneMapCheckResult::SkipScan {
820 continue;
821 }
822 }
823
824 for row in 0..group.num_nodes as usize {
825 ids.push(group.start_offset + row as u64);
826 for (col, res_col) in result.iter_mut().enumerate().take(num_cols) {
827 match group.get_value_owned_with_snapshot(row, col, snapshot_ts, commit_history) {
828 Some(v) => res_col.push(v),
829 None => res_col.push(Value::Null),
830 }
831 }
832 }
833 }
834
835 (result, ids)
836 }
837
838 fn find_group(&self, row: u64) -> usize {
840 self.node_groups
841 .binary_search_by_key(&row, |g| g.start_offset)
842 .unwrap_or_else(|i| if i == 0 { 0 } else { i - 1 })
843 }
844}
845
846fn coerce_values_to_columns(values: &mut [Value], columns: &[ColumnDefinition]) -> Result<(), StorageError> {
854 for (i, col) in columns.iter().enumerate() {
855 let coerced = match (col.logical_type, &values[i]) {
856 (LogicalTypeID::UInt64, Value::Int64(x)) if *x >= 0 => Some(Value::UInt64(*x as u64)),
857 (LogicalTypeID::UInt64, Value::Int32(x)) if *x >= 0 => Some(Value::UInt64(*x as u64)),
858 (LogicalTypeID::UInt64, Value::Int16(x)) if *x >= 0 => Some(Value::UInt64(*x as u64)),
859 (LogicalTypeID::UInt64, Value::Int8(x)) if *x >= 0 => Some(Value::UInt64(*x as u64)),
860 (LogicalTypeID::UInt64, Value::UInt32(x)) => Some(Value::UInt64(*x as u64)),
861 (LogicalTypeID::UInt64, Value::UInt16(x)) => Some(Value::UInt64(*x as u64)),
862 (LogicalTypeID::UInt64, Value::UInt8(x)) => Some(Value::UInt64(*x as u64)),
863 (LogicalTypeID::UInt64, Value::Int64(_))
864 | (LogicalTypeID::UInt64, Value::Int32(_))
865 | (LogicalTypeID::UInt64, Value::Int16(_))
866 | (LogicalTypeID::UInt64, Value::Int8(_)) => {
867 return Err(StorageError::Page(format!(
868 "Cannot store negative value in UINT64 column '{}'",
869 col.name
870 )));
871 }
872 _ => None,
873 };
874 if let Some(c) = coerced {
875 values[i] = c;
876 }
877 }
878 Ok(())
879}
880
881#[derive(Debug, Clone)]
893pub struct RelTable {
894 pub table_id: u64,
895 pub name: String,
896 pub src_table_id: u64,
897 pub dst_table_id: u64,
898 pub columns: Vec<ColumnDefinition>,
899 pub num_rows: u64,
900 pub edges: Vec<(u64, u64)>,
902 pub fwd_adj: HashMap<u64, Vec<(u64, usize)>>,
904 pub rev_adj: HashMap<u64, Vec<(u64, usize)>>,
906 pub csr_index: Option<CsrIndex>,
908 pub properties: Vec<Vec<Value>>,
910 pub persistence_dirty: bool,
913}
914
915impl RelTable {
916 pub fn new(
917 table_id: u64,
918 name: String,
919 src_table_id: u64,
920 dst_table_id: u64,
921 columns: Vec<ColumnDefinition>,
922 ) -> Self {
923 let num_cols = columns.len();
924 Self {
925 table_id,
926 name,
927 src_table_id,
928 dst_table_id,
929 columns,
930 num_rows: 0,
931 edges: Vec::new(),
932 fwd_adj: HashMap::new(),
933 rev_adj: HashMap::new(),
934 csr_index: None,
935 properties: vec![Vec::new(); num_cols],
936 persistence_dirty: false,
937 }
938 }
939
940 pub fn add_column(&mut self, column: ColumnDefinition) {
943 if self.columns.iter().any(|c| c.name.eq_ignore_ascii_case(&column.name)) {
944 return;
945 }
946 self.columns.push(column);
947 self.properties.push(vec![Value::Null; self.edges.len()]);
948 self.persistence_dirty = true;
949 }
950
951 pub fn insert_rel(&mut self, from: u64, to: u64, values: Vec<Value>) -> Result<(), StorageError> {
959 if values.len() != self.columns.len() {
960 return Err(StorageError::Page(format!(
961 "Column count mismatch: expected {} values, got {}",
962 self.columns.len(),
963 values.len()
964 )));
965 }
966
967 let edge_idx = self.edges.len();
968 self.edges.push((from, to));
969
970 self.fwd_adj.entry(from).or_default().push((to, edge_idx));
972
973 self.rev_adj.entry(to).or_default().push((from, edge_idx));
975 for (col_idx, val) in values.into_iter().enumerate() {
977 self.properties[col_idx].push(val);
978 }
979 self.num_rows += 1;
980 self.persistence_dirty = true;
983 Ok(())
984 }
985
986 pub fn insert_rels_batch(&mut self, rels: &[(u64, u64, Vec<Value>)]) -> Result<u64, StorageError> {
989 if rels.is_empty() {
990 return Ok(0);
991 }
992 let num_cols = self.columns.len();
993 let total = rels.len();
994
995 for (i, (_, _, vals)) in rels.iter().enumerate() {
997 if vals.len() != num_cols {
998 return Err(StorageError::Page(format!(
999 "Rel {} column count mismatch: expected {} values, got {}",
1000 i,
1001 num_cols,
1002 vals.len()
1003 )));
1004 }
1005 }
1006
1007 self.edges.reserve(total);
1009 for col in &mut self.properties {
1010 col.reserve(total);
1011 }
1012
1013 let _start_edge_idx = self.edges.len();
1014
1015 for (from, to, vals) in rels {
1017 let edge_idx = self.edges.len();
1018 self.edges.push((*from, *to));
1019 self.fwd_adj.entry(*from).or_default().push((*to, edge_idx));
1020 self.rev_adj.entry(*to).or_default().push((*from, edge_idx));
1021 for (col_idx, val) in vals.iter().enumerate() {
1022 self.properties[col_idx].push(val.clone());
1023 }
1024 }
1025
1026 self.num_rows += total as u64;
1027 self.persistence_dirty = true;
1030 Ok(total as u64)
1031 }
1032
1033 pub fn delete_edge(&mut self, edge_idx: usize) -> Result<(), StorageError> {
1036 if edge_idx >= self.edges.len() {
1037 return Err(StorageError::Page(format!("Edge index {edge_idx} out of range")));
1038 }
1039
1040 let (src, dst) = self.edges[edge_idx];
1041 if src == u64::MAX {
1042 return Ok(());
1044 }
1045
1046 if let Some(adj) = self.fwd_adj.get_mut(&src) {
1048 adj.retain(|&(_, idx)| idx != edge_idx);
1049 }
1050
1051 if let Some(adj) = self.rev_adj.get_mut(&dst) {
1053 adj.retain(|&(_, idx)| idx != edge_idx);
1054 }
1055
1056 self.edges[edge_idx] = (u64::MAX, u64::MAX);
1058 self.persistence_dirty = true;
1059
1060 for col in &mut self.properties {
1062 if edge_idx < col.len() {
1063 col[edge_idx] = Value::Null;
1064 }
1065 }
1066
1067 Ok(())
1068 }
1069
1070 pub fn update_cell(&mut self, edge_idx: usize, col_idx: usize, value: Value) -> Result<(), StorageError> {
1072 if col_idx >= self.columns.len() {
1073 return Err(StorageError::Page(format!("Column index {col_idx} out of range")));
1074 }
1075 if edge_idx >= self.properties[col_idx].len() {
1076 return Err(StorageError::Page(format!("Edge index {edge_idx} out of range")));
1077 }
1078
1079 self.properties[col_idx][edge_idx] = value;
1080 self.persistence_dirty = true;
1081 Ok(())
1082 }
1083
1084 pub fn edge_undo_bytes(&self, edge_idx: usize) -> Vec<u8> {
1088 let (src, dst) = self.edges.get(edge_idx).copied().unwrap_or((u64::MAX, u64::MAX));
1089 let mut out = Vec::new();
1090 out.extend_from_slice(&Column::serialize_value(&Value::UInt64(src)));
1091 out.extend_from_slice(&Column::serialize_value(&Value::UInt64(dst)));
1092 for p in self.get_edge_properties(edge_idx) {
1093 out.extend_from_slice(&Column::serialize_value(&p));
1094 }
1095 out
1096 }
1097
1098 pub fn edge_cell_undo_bytes(&self, edge_idx: usize, col_idx: usize) -> Vec<u8> {
1101 let v = self
1102 .get_edge_properties(edge_idx)
1103 .get(col_idx)
1104 .cloned()
1105 .unwrap_or(Value::Null);
1106 Column::serialize_value(&v)
1107 }
1108
1109 pub fn restore_deleted_edge(
1113 &mut self,
1114 edge_idx: usize,
1115 src: u64,
1116 dst: u64,
1117 props: Vec<Value>,
1118 ) -> Result<(), StorageError> {
1119 if edge_idx >= self.edges.len() {
1120 return Err(StorageError::Page(format!("Edge index {edge_idx} out of range")));
1121 }
1122 self.edges[edge_idx] = (src, dst);
1123 self.fwd_adj.entry(src).or_default().push((dst, edge_idx));
1124 self.rev_adj.entry(dst).or_default().push((src, edge_idx));
1125 for (col_idx, val) in props.into_iter().enumerate() {
1126 if col_idx < self.properties.len() {
1127 if edge_idx < self.properties[col_idx].len() {
1128 self.properties[col_idx][edge_idx] = val;
1129 } else {
1130 self.properties[col_idx].push(val);
1131 }
1132 }
1133 }
1134 self.persistence_dirty = true;
1135 Ok(())
1136 }
1137
1138 pub fn insert_row(&mut self, values: Vec<Value>) -> Result<u64, StorageError> {
1142 let num_prop_cols = self.columns.len();
1146 if values.len() != num_prop_cols {
1147 return Err(StorageError::Page(format!(
1148 "Column count mismatch: expected {} values, got {}",
1149 num_prop_cols,
1150 values.len()
1151 )));
1152 }
1153
1154 let from = self.num_rows;
1157 let to = self.num_rows;
1158 self.insert_rel(from, to, values)?;
1159 Ok(0) }
1161
1162 pub fn scan_adj_list(&self, src_offset: u64) -> &[(u64, usize)] {
1167 self.fwd_adj.get(&src_offset).map(|v| v.as_slice()).unwrap_or(&[])
1168 }
1169
1170 pub fn scan_rev_adj_list(&self, dst_offset: u64) -> &[(u64, usize)] {
1175 self.rev_adj.get(&dst_offset).map(|v| v.as_slice()).unwrap_or(&[])
1176 }
1177
1178 pub fn get_outgoing_edges(&self, src_offset: u64) -> Vec<(u64, Vec<Value>)> {
1180 self.scan_adj_list(src_offset)
1181 .iter()
1182 .map(|&(dst, edge_idx)| {
1183 let props = self.get_edge_properties(edge_idx);
1184 (dst, props)
1185 })
1186 .collect()
1187 }
1188
1189 pub fn get_incoming_edges(&self, dst_offset: u64) -> Vec<(u64, Vec<Value>)> {
1191 self.scan_rev_adj_list(dst_offset)
1192 .iter()
1193 .map(|&(src, edge_idx)| {
1194 let props = self.get_edge_properties(edge_idx);
1195 (src, props)
1196 })
1197 .collect()
1198 }
1199
1200 pub fn get_edge_properties(&self, edge_idx: usize) -> Vec<Value> {
1202 let mut props = Vec::with_capacity(self.columns.len());
1203 for col in &self.properties {
1204 match col.get(edge_idx) {
1205 Some(v) => props.push(v.clone()),
1206 None => props.push(Value::Null),
1207 }
1208 }
1209 props
1210 }
1211
1212 pub fn get_column(&self, col_idx: usize) -> Option<&[Value]> {
1214 self.properties.get(col_idx).map(|v| v.as_slice())
1215 }
1216
1217 pub fn to_column_major_data(&self) -> Vec<Vec<Value>> {
1219 self.properties.clone()
1220 }
1221}
1222
1223#[derive(Debug, Default)]
1230pub struct TableCatalog {
1231 node_tables: DashMap<u64, NodeTable>,
1232 rel_tables: DashMap<u64, RelTable>,
1233 vector_indexes: DashMap<u64, VectorIndexTable>,
1234 node_name_to_id: DashMap<String, u64>,
1236 rel_name_to_id: DashMap<String, u64>,
1238 vector_index_name_to_id: DashMap<String, u64>,
1240 next_table_id: std::sync::atomic::AtomicU64,
1241}
1242
1243impl TableCatalog {
1244 pub fn new() -> Self {
1245 Self::default()
1246 }
1247
1248 pub fn create_node_table(&self, name: String, columns: Vec<ColumnDefinition>) -> NodeTable {
1249 let table_id = self.next_table_id.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1250 let table = NodeTable::new(table_id, name.clone(), columns);
1251 self.node_name_to_id.insert(name, table_id);
1252 self.node_tables.insert(table_id, table.clone());
1253 table
1254 }
1255
1256 pub fn create_node_table_with_id(&self, table_id: u64, name: String, columns: Vec<ColumnDefinition>) -> NodeTable {
1260 self.bump_next_table_id(table_id);
1261 let table = NodeTable::new(table_id, name.clone(), columns);
1262 self.node_name_to_id.insert(name, table_id);
1263 self.node_tables.insert(table_id, table.clone());
1264 table
1265 }
1266
1267 pub fn create_rel_table(
1268 &self,
1269 name: String,
1270 src_table_id: u64,
1271 dst_table_id: u64,
1272 columns: Vec<ColumnDefinition>,
1273 ) -> RelTable {
1274 let table_id = self.next_table_id.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1275 let table = RelTable::new(table_id, name.clone(), src_table_id, dst_table_id, columns);
1276 self.rel_name_to_id.insert(name, table_id);
1277 self.rel_tables.insert(table_id, table.clone());
1278 table
1279 }
1280
1281 pub fn create_rel_table_with_id(
1285 &self,
1286 table_id: u64,
1287 name: String,
1288 src_table_id: u64,
1289 dst_table_id: u64,
1290 columns: Vec<ColumnDefinition>,
1291 ) -> RelTable {
1292 self.bump_next_table_id(table_id);
1293 let table = RelTable::new(table_id, name.clone(), src_table_id, dst_table_id, columns);
1294 self.rel_name_to_id.insert(name, table_id);
1295 self.rel_tables.insert(table_id, table.clone());
1296 table
1297 }
1298
1299 fn bump_next_table_id(&self, table_id: u64) {
1302 let mut next = self.next_table_id.load(std::sync::atomic::Ordering::SeqCst);
1303 while next <= table_id {
1304 match self.next_table_id.compare_exchange(
1305 next,
1306 table_id + 1,
1307 std::sync::atomic::Ordering::SeqCst,
1308 std::sync::atomic::Ordering::SeqCst,
1309 ) {
1310 Ok(_) => break,
1311 Err(current) => next = current,
1312 }
1313 }
1314 }
1315
1316 pub fn get_node_table(&self, table_id: u64) -> Option<dashmap::mapref::one::Ref<'_, u64, NodeTable>> {
1317 self.node_tables.get(&table_id)
1318 }
1319
1320 pub fn get_node_table_mut(&self, table_id: u64) -> Option<dashmap::mapref::one::RefMut<'_, u64, NodeTable>> {
1321 self.node_tables.get_mut(&table_id)
1322 }
1323
1324 pub fn get_node_table_by_name(&self, name: &str) -> Option<dashmap::mapref::one::Ref<'_, u64, NodeTable>> {
1325 let id = self.node_name_to_id.get(name)?;
1326 self.node_tables.get(&*id)
1327 }
1328
1329 pub fn get_node_table_by_name_mut(&self, name: &str) -> Option<dashmap::mapref::one::RefMut<'_, u64, NodeTable>> {
1330 let id = self.node_name_to_id.get(name)?;
1331 self.node_tables.get_mut(&*id)
1332 }
1333
1334 pub fn get_rel_table(&self, table_id: u64) -> Option<dashmap::mapref::one::Ref<'_, u64, RelTable>> {
1335 self.rel_tables.get(&table_id)
1336 }
1337
1338 pub fn get_rel_table_mut(&self, table_id: u64) -> Option<dashmap::mapref::one::RefMut<'_, u64, RelTable>> {
1339 self.rel_tables.get_mut(&table_id)
1340 }
1341
1342 pub fn get_rel_table_by_name(&self, name: &str) -> Option<dashmap::mapref::one::Ref<'_, u64, RelTable>> {
1343 let id = self.rel_name_to_id.get(name)?;
1344 self.rel_tables.get(&*id)
1345 }
1346
1347 pub fn get_rel_table_by_name_mut(&self, name: &str) -> Option<dashmap::mapref::one::RefMut<'_, u64, RelTable>> {
1348 let id = self.rel_name_to_id.get(name)?;
1349 self.rel_tables.get_mut(&*id)
1350 }
1351
1352 pub fn has_incident_edges(&self, table_id: u64, node_idx: u64) -> bool {
1354 for rel_table in self.rel_tables.iter() {
1355 if rel_table.src_table_id == table_id {
1356 if let Some(edges) = rel_table.fwd_adj.get(&node_idx) {
1357 if !edges.is_empty() {
1358 return true;
1359 }
1360 }
1361 }
1362 if rel_table.dst_table_id == table_id {
1363 if let Some(edges) = rel_table.rev_adj.get(&node_idx) {
1364 if !edges.is_empty() {
1365 return true;
1366 }
1367 }
1368 }
1369 }
1370 false
1371 }
1372
1373 pub fn detach_node(&self, table_id: u64, node_idx: u64) {
1375 for mut rel_table in self.rel_tables.iter_mut() {
1376 let mut edges_to_delete = Vec::new();
1377
1378 if rel_table.src_table_id == table_id {
1379 if let Some(edges) = rel_table.fwd_adj.get(&node_idx) {
1380 for &(_, edge_idx) in edges {
1381 edges_to_delete.push(edge_idx);
1382 }
1383 }
1384 }
1385
1386 if rel_table.dst_table_id == table_id {
1387 if let Some(edges) = rel_table.rev_adj.get(&node_idx) {
1388 for &(_, edge_idx) in edges {
1389 edges_to_delete.push(edge_idx);
1390 }
1391 }
1392 }
1393
1394 for edge_idx in edges_to_delete {
1395 let _ = rel_table.delete_edge(edge_idx);
1396 }
1397 }
1398 }
1399
1400 pub fn all_node_tables(&self) -> Vec<dashmap::mapref::multiple::RefMulti<'_, u64, NodeTable>> {
1401 self.node_tables.iter().collect()
1402 }
1403
1404 pub fn all_rel_tables(&self) -> Vec<dashmap::mapref::multiple::RefMulti<'_, u64, RelTable>> {
1405 self.rel_tables.iter().collect()
1406 }
1407
1408 pub fn node_table_num_rows(&self, name: &str) -> u64 {
1410 self.get_node_table_by_name(name).map(|t| t.num_rows).unwrap_or(0)
1411 }
1412
1413 pub fn drop_node_table(&self, name: &str) -> bool {
1415 if let Some(id) = self.node_name_to_id.get(name) {
1416 let table_id = *id;
1417 drop(id);
1418 self.node_name_to_id.remove(name);
1419 self.node_tables.remove(&table_id).is_some()
1420 } else {
1421 false
1422 }
1423 }
1424
1425 pub fn create_vector_index(
1427 &self,
1428 name: String,
1429 table_name: String,
1430 column_name: String,
1431 metric: DistanceMetric,
1432 dimensions: u32,
1433 ) -> VectorIndexTable {
1434 let index_id = self.next_table_id.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1435 let table = VectorIndexTable::new(index_id, name.clone(), table_name, column_name, metric, dimensions);
1436 self.vector_index_name_to_id.insert(name, index_id);
1437 self.vector_indexes.insert(index_id, table.clone());
1438 table
1439 }
1440
1441 pub fn get_vector_index(&self, index_id: u64) -> Option<dashmap::mapref::one::Ref<'_, u64, VectorIndexTable>> {
1443 self.vector_indexes.get(&index_id)
1444 }
1445
1446 pub fn get_vector_index_by_name(&self, name: &str) -> Option<dashmap::mapref::one::Ref<'_, u64, VectorIndexTable>> {
1448 let id = self.vector_index_name_to_id.get(name)?;
1449 self.vector_indexes.get(&*id)
1450 }
1451
1452 pub fn get_vector_index_by_name_mut(
1454 &self,
1455 name: &str,
1456 ) -> Option<dashmap::mapref::one::RefMut<'_, u64, VectorIndexTable>> {
1457 let id = self.vector_index_name_to_id.get(name)?;
1458 self.vector_indexes.get_mut(&*id)
1459 }
1460
1461 pub fn get_vector_index_mut(
1463 &self,
1464 index_id: u64,
1465 ) -> Option<dashmap::mapref::one::RefMut<'_, u64, VectorIndexTable>> {
1466 self.vector_indexes.get_mut(&index_id)
1467 }
1468
1469 pub fn drop_vector_index(&self, name: &str) -> bool {
1471 if let Some(id) = self.vector_index_name_to_id.get(name) {
1472 let index_id = *id;
1473 drop(id);
1474 self.vector_index_name_to_id.remove(name);
1475 self.vector_indexes.remove(&index_id).is_some()
1476 } else {
1477 false
1478 }
1479 }
1480
1481 pub fn all_vector_indexes(&self) -> Vec<dashmap::mapref::multiple::RefMulti<'_, u64, VectorIndexTable>> {
1483 self.vector_indexes.iter().collect()
1484 }
1485
1486 pub fn refresh_vector_index(&self, index_id: u64) {
1491 let (table_name, column_name) = match self.vector_indexes.get(&index_id) {
1492 Some(vi) => (vi.table_name.clone(), vi.column_name.clone()),
1493 None => return,
1494 };
1495
1496 let col_idx = match self.get_node_table_by_name(&table_name) {
1497 Some(t) => match t.columns.iter().position(|c| c.name == column_name) {
1498 Some(idx) => idx,
1499 None => return,
1500 },
1501 None => return,
1502 };
1503
1504 let mut data: Vec<(usize, Vec<f64>)> = Vec::new();
1507 if let Some(table) = self.get_node_table_by_name(&table_name) {
1508 for row_id in 0..table.num_rows as usize {
1509 if let Some(val) = table.get_value(row_id, col_idx) {
1510 if let Ok(vec) = crate::extract_f64_list_from_value(val) {
1511 data.push((row_id, vec));
1512 }
1513 }
1514 }
1515 }
1516 if data.is_empty() {
1517 if let Some(mut vi) = self.vector_indexes.get_mut(&index_id) {
1518 vi.hnsw_mut().clear();
1519 }
1520 return;
1521 }
1522
1523 let mut vi = self.vector_indexes.get_mut(&index_id);
1524 if let Some(vi) = vi.as_mut() {
1525 vi.hnsw_mut().clear();
1526 for (row_id, vec) in data {
1527 vi.hnsw_mut().insert(vec, row_id);
1528 }
1529 }
1530 }
1531
1532 pub fn refresh_vector_indexes_for_tables(&self, table_ids: &[u64]) {
1537 for table_id in table_ids {
1538 let table_name = match self.get_node_table(*table_id) {
1539 Some(t) => t.name.clone(),
1540 None => continue,
1541 };
1542 let index_ids: Vec<u64> = self
1543 .vector_indexes
1544 .iter()
1545 .filter(|vi| vi.table_name == table_name)
1546 .map(|vi| *vi.key())
1547 .collect();
1548 for index_id in index_ids {
1549 self.refresh_vector_index(index_id);
1550 }
1551 }
1552 }
1553
1554 pub fn create_art_index(&self, table_name: &str, index_name: &str) -> Result<(), StorageError> {
1561 let mut table = self
1562 .get_node_table_by_name_mut(table_name)
1563 .ok_or_else(|| StorageError::TableNotFound(format!("Node table '{table_name}' not found")))?;
1564
1565 if table.art_index.is_some() {
1566 return Err(StorageError::Index(format!(
1567 "Table '{table_name}' already has an ART index"
1568 )));
1569 }
1570
1571 let mut art_idx = ArtPrimaryKeyIndex::new(index_name);
1572
1573 let pk_col = table.primary_key_column;
1575 let col_major = table.to_column_major_data();
1577 if pk_col < col_major.len() {
1578 for (row_offset, pk_val) in col_major[pk_col].iter().enumerate() {
1579 if !matches!(pk_val, Value::Null)
1580 && let Some(art_key) = ArtKey::from_value(pk_val)
1581 {
1582 art_idx.insert(&art_key, row_offset as u64);
1583 }
1584 }
1585 }
1586
1587 table.art_index = Some(art_idx);
1588 Ok(())
1589 }
1590
1591 pub fn drop_art_index(&self, table_name: &str) -> Result<(), StorageError> {
1593 let mut table = self
1594 .get_node_table_by_name_mut(table_name)
1595 .ok_or_else(|| StorageError::TableNotFound(format!("Node table '{table_name}' not found")))?;
1596
1597 table.art_index = None;
1598 Ok(())
1599 }
1600
1601 pub fn get_art_index(&self, table_name: &str) -> Option<ArtPrimaryKeyIndex> {
1604 let table = self.get_node_table_by_name(table_name)?;
1605 table.art_index.clone()
1606 }
1607
1608 pub fn has_art_index(&self, table_name: &str) -> bool {
1610 self.get_node_table_by_name(table_name)
1611 .map(|t| t.art_index.is_some())
1612 .unwrap_or(false)
1613 }
1614
1615 pub fn drop_rel_table(&self, name: &str) -> bool {
1617 if let Some(id) = self.rel_name_to_id.get(name) {
1618 let table_id = *id;
1619 drop(id);
1620 self.rel_name_to_id.remove(name);
1621 self.rel_tables.remove(&table_id).is_some()
1622 } else {
1623 false
1624 }
1625 }
1626}
1627
1628#[cfg(test)]
1633mod tests {
1634 use super::*;
1635 use crate::column_chunk::NODE_GROUP_SIZE;
1636 use std::sync::Arc;
1637
1638 #[test]
1641 fn test_node_table_empty() {
1642 let table = NodeTable::new(
1643 1,
1644 "Person".into(),
1645 vec![
1646 ColumnDefinition {
1647 compression: akar_common::enums::CompressionType::Uncompressed,
1648 name: "name".into(),
1649 logical_type: LogicalTypeID::String,
1650 is_primary_key: true,
1651 },
1652 ColumnDefinition {
1653 compression: akar_common::enums::CompressionType::Uncompressed,
1654 name: "age".into(),
1655 logical_type: LogicalTypeID::Int64,
1656 is_primary_key: false,
1657 },
1658 ],
1659 );
1660 assert_eq!(table.num_rows, 0);
1661 assert!(table.node_groups.is_empty());
1662 }
1663
1664 #[test]
1665 fn test_node_table_insert_and_get() {
1666 let mut table = NodeTable::new(
1667 1,
1668 "Person".into(),
1669 vec![
1670 ColumnDefinition {
1671 compression: akar_common::enums::CompressionType::Uncompressed,
1672 name: "name".into(),
1673 logical_type: LogicalTypeID::String,
1674 is_primary_key: true,
1675 },
1676 ColumnDefinition {
1677 compression: akar_common::enums::CompressionType::Uncompressed,
1678 name: "age".into(),
1679 logical_type: LogicalTypeID::Int64,
1680 is_primary_key: false,
1681 },
1682 ],
1683 );
1684 table
1685 .insert_row(vec![Value::String("Alice".into()), Value::Int64(30)])
1686 .unwrap();
1687 table
1688 .insert_row(vec![Value::String("Bob".into()), Value::Int64(25)])
1689 .unwrap();
1690
1691 assert_eq!(table.num_rows, 2);
1692 assert_eq!(table.get_value(0, 0), Some(&Value::String("Alice".into())));
1693 assert_eq!(table.get_value(1, 1), Some(&Value::Int64(25)));
1694 }
1695
1696 #[test]
1697 fn test_node_table_batch_insert_with_spiller_restores_all_rows() {
1698 let dir = tempfile::tempdir().unwrap();
1699 let spiller = Arc::new(crate::spiller::Spiller::new(dir.path(), 64));
1700 let mut table = NodeTable::new(
1701 1,
1702 "T".into(),
1703 vec![ColumnDefinition {
1704 compression: akar_common::enums::CompressionType::Uncompressed,
1705 name: "id".into(),
1706 logical_type: LogicalTypeID::Int64,
1707 is_primary_key: true,
1708 }],
1709 );
1710 table.set_spiller(Some(spiller));
1711
1712 let rows: Vec<Vec<Value>> = (0..2000).map(|i| vec![Value::Int64(i)]).collect();
1713 let inserted = table.insert_rows_batch(&rows).unwrap();
1714 assert_eq!(inserted, 2000);
1715 assert_eq!(table.num_rows, 2000);
1716
1717 for i in 0i64..2000 {
1718 assert_eq!(table.get_value(i as usize, 0), Some(&Value::Int64(i)));
1719 }
1720 assert!(
1721 table.node_groups.iter().all(|g| !g.has_spill_files()),
1722 "spill files must be merged back after the batch"
1723 );
1724 }
1725
1726 #[test]
1727 fn test_delete_row_removes_pk_from_hash_and_art_index() {
1728 let mut table = NodeTable::new(
1729 1,
1730 "Person".into(),
1731 vec![
1732 ColumnDefinition {
1733 compression: akar_common::enums::CompressionType::Uncompressed,
1734 name: "name".into(),
1735 logical_type: LogicalTypeID::String,
1736 is_primary_key: true,
1737 },
1738 ColumnDefinition {
1739 compression: akar_common::enums::CompressionType::Uncompressed,
1740 name: "age".into(),
1741 logical_type: LogicalTypeID::Int64,
1742 is_primary_key: false,
1743 },
1744 ],
1745 );
1746 table.art_index = Some(ArtPrimaryKeyIndex::new("test_art"));
1747 table
1748 .insert_row(vec![Value::String("Alice".into()), Value::Int64(30)])
1749 .unwrap();
1750 table
1751 .insert_row(vec![Value::String("Bob".into()), Value::Int64(25)])
1752 .unwrap();
1753 let art = table.art_index.as_ref().unwrap();
1754 assert_eq!(
1755 art.lookup(&ArtKey::from_value(&Value::String("Alice".into())).unwrap()),
1756 Some(0)
1757 );
1758
1759 table.delete_row(0).unwrap();
1760
1761 assert!(table.lookup_by_pk(&Value::String("Alice".into())).is_none());
1763 let art = table.art_index.as_ref().unwrap();
1764 assert_eq!(art.len(), 1, "ART must drop the deleted entry");
1765 assert!(
1766 art.lookup(&ArtKey::from_value(&Value::String("Alice".into())).unwrap())
1767 .is_none()
1768 );
1769 assert!(
1770 art.lookup(&ArtKey::from_value(&Value::String("Bob".into())).unwrap())
1771 .is_some()
1772 );
1773 let hits = table.lookup_by_pk_range(
1775 Some(&Value::String("A".into())),
1776 true,
1777 Some(&Value::String("C".into())),
1778 true,
1779 100,
1780 );
1781 assert_eq!(hits, vec![1], "only 'Bob' (row 1) should be in range");
1782
1783 table
1785 .insert_row(vec![Value::String("Alice".into()), Value::Int64(31)])
1786 .unwrap();
1787 assert_eq!(table.lookup_by_pk(&Value::String("Alice".into())), Some(2));
1788 }
1789
1790 #[test]
1791 fn test_node_table_scan_column() {
1792 let mut table = NodeTable::new(
1793 1,
1794 "T".into(),
1795 vec![ColumnDefinition {
1796 compression: akar_common::enums::CompressionType::Uncompressed,
1797 name: "val".into(),
1798 logical_type: LogicalTypeID::Int64,
1799 is_primary_key: false,
1800 }],
1801 );
1802 for i in 0..100 {
1803 table.insert_row(vec![Value::Int64(i)]).unwrap();
1804 }
1805 let scanned = table.scan_column(0, 10, 5, None, &HashMap::new());
1806 assert_eq!(scanned.len(), 5);
1807 assert_eq!(scanned[0], Value::Int64(10));
1808 assert_eq!(scanned[4], Value::Int64(14));
1809 }
1810
1811 #[test]
1812 fn test_node_table_to_column_major() {
1813 let mut table = NodeTable::new(
1814 1,
1815 "T".into(),
1816 vec![
1817 ColumnDefinition {
1818 compression: akar_common::enums::CompressionType::Uncompressed,
1819 name: "x".into(),
1820 logical_type: LogicalTypeID::Int64,
1821 is_primary_key: false,
1822 },
1823 ColumnDefinition {
1824 compression: akar_common::enums::CompressionType::Uncompressed,
1825 name: "y".into(),
1826 logical_type: LogicalTypeID::Int64,
1827 is_primary_key: false,
1828 },
1829 ],
1830 );
1831 table.insert_row(vec![Value::Int64(1), Value::Int64(10)]).unwrap();
1832 table.insert_row(vec![Value::Int64(2), Value::Int64(20)]).unwrap();
1833
1834 let data = table.to_column_major_data();
1835 assert_eq!(data.len(), 2);
1836 assert_eq!(data[0], vec![Value::Int64(1), Value::Int64(2)]);
1837 assert_eq!(data[1], vec![Value::Int64(10), Value::Int64(20)]);
1838 }
1839
1840 #[test]
1841 fn test_node_table_auto_node_group() {
1842 let mut table = NodeTable::new(
1843 1,
1844 "T".into(),
1845 vec![ColumnDefinition {
1846 compression: akar_common::enums::CompressionType::Uncompressed,
1847 name: "v".into(),
1848 logical_type: LogicalTypeID::Int64,
1849 is_primary_key: false,
1850 }],
1851 );
1852 for i in 0..NODE_GROUP_SIZE as u64 + 1 {
1854 table.insert_row(vec![Value::Int64(i as i64)]).unwrap();
1855 }
1856 assert_eq!(table.num_rows, NODE_GROUP_SIZE as u64 + 1);
1857 assert_eq!(table.node_groups.len(), 2);
1858 assert_eq!(table.node_groups[0].num_nodes, NODE_GROUP_SIZE as u64);
1859 assert_eq!(table.node_groups[1].num_nodes, 1);
1860 assert_eq!(table.get_value(0, 0), Some(&Value::Int64(0)));
1862 assert_eq!(
1863 table.get_value(NODE_GROUP_SIZE, 0),
1864 Some(&Value::Int64(NODE_GROUP_SIZE as i64))
1865 );
1866 }
1867
1868 fn make_rel_table() -> RelTable {
1871 RelTable::new(
1872 1,
1873 "Knows".into(),
1874 0,
1875 1,
1876 vec![
1877 ColumnDefinition {
1878 compression: akar_common::enums::CompressionType::Uncompressed,
1879 name: "since".into(),
1880 logical_type: LogicalTypeID::Int64,
1881 is_primary_key: false,
1882 },
1883 ColumnDefinition {
1884 compression: akar_common::enums::CompressionType::Uncompressed,
1885 name: "weight".into(),
1886 logical_type: LogicalTypeID::Double,
1887 is_primary_key: false,
1888 },
1889 ],
1890 )
1891 }
1892
1893 #[test]
1894 fn test_rel_table_empty() {
1895 let rel = make_rel_table();
1896 assert_eq!(rel.num_rows, 0);
1897 assert!(rel.edges.is_empty());
1898 assert!(rel.fwd_adj.is_empty());
1899 assert!(rel.rev_adj.is_empty());
1900 }
1901
1902 #[test]
1903 fn test_rel_insert_basic() {
1904 let mut rel = make_rel_table();
1905 rel.insert_rel(0, 1, vec![Value::Int64(2020), Value::Double(0.5)])
1906 .unwrap();
1907 rel.insert_rel(0, 2, vec![Value::Int64(2021), Value::Double(0.8)])
1908 .unwrap();
1909 rel.insert_rel(1, 0, vec![Value::Int64(2020), Value::Double(0.3)])
1910 .unwrap();
1911
1912 assert_eq!(rel.num_rows, 3);
1913 assert_eq!(rel.edges.len(), 3);
1914
1915 let fwd = rel.scan_adj_list(0);
1917 assert_eq!(fwd.len(), 2);
1918 assert_eq!(fwd[0], (1, 0)); assert_eq!(fwd[1], (2, 1)); let fwd1 = rel.scan_adj_list(1);
1923 assert_eq!(fwd1.len(), 1);
1924 assert_eq!(fwd1[0], (0, 2));
1925 }
1926
1927 #[test]
1928 fn test_rel_reverse_adjacency() {
1929 let mut rel = make_rel_table();
1930 rel.insert_rel(0, 5, vec![Value::Int64(2022), Value::Double(1.0)])
1931 .unwrap();
1932 rel.insert_rel(3, 5, vec![Value::Int64(2023), Value::Double(1.5)])
1933 .unwrap();
1934
1935 let rev = rel.scan_rev_adj_list(5);
1937 assert_eq!(rev.len(), 2);
1938 assert_eq!(rev[0], (0, 0));
1939 assert_eq!(rev[1], (3, 1));
1940 }
1941
1942 #[test]
1943 fn test_rel_get_edge_properties() {
1944 let mut rel = make_rel_table();
1945 rel.insert_rel(0, 1, vec![Value::Int64(2020), Value::Double(0.5)])
1946 .unwrap();
1947 rel.insert_rel(2, 3, vec![Value::Int64(2021), Value::Double(0.9)])
1948 .unwrap();
1949
1950 let props0 = rel.get_edge_properties(0);
1951 assert_eq!(props0, vec![Value::Int64(2020), Value::Double(0.5)]);
1952
1953 let props1 = rel.get_edge_properties(1);
1954 assert_eq!(props1, vec![Value::Int64(2021), Value::Double(0.9)]);
1955 }
1956
1957 #[test]
1958 fn test_rel_get_outgoing_edges() {
1959 let mut rel = make_rel_table();
1960 rel.insert_rel(0, 10, vec![Value::Int64(2020), Value::Double(1.0)])
1961 .unwrap();
1962 rel.insert_rel(0, 20, vec![Value::Int64(2021), Value::Double(2.0)])
1963 .unwrap();
1964
1965 let outgoing = rel.get_outgoing_edges(0);
1966 assert_eq!(outgoing.len(), 2);
1967 assert_eq!(outgoing[0].0, 10);
1968 assert_eq!(outgoing[0].1, vec![Value::Int64(2020), Value::Double(1.0)]);
1969 assert_eq!(outgoing[1].0, 20);
1970 }
1971
1972 #[test]
1973 fn test_rel_get_incoming_edges() {
1974 let mut rel = make_rel_table();
1975 rel.insert_rel(10, 5, vec![Value::Int64(2020), Value::Double(1.0)])
1976 .unwrap();
1977 rel.insert_rel(20, 5, vec![Value::Int64(2021), Value::Double(2.0)])
1978 .unwrap();
1979
1980 let incoming = rel.get_incoming_edges(5);
1981 assert_eq!(incoming.len(), 2);
1982 assert_eq!(incoming[0].0, 10);
1983 assert_eq!(incoming[1].0, 20);
1984 }
1985
1986 #[test]
1987 fn test_rel_no_edges() {
1988 let rel = make_rel_table();
1989 assert!(rel.scan_adj_list(0).is_empty());
1990 assert!(rel.scan_rev_adj_list(0).is_empty());
1991 assert!(rel.get_outgoing_edges(0).is_empty());
1992 assert!(rel.get_incoming_edges(0).is_empty());
1993 }
1994
1995 #[test]
1996 fn test_rel_insert_row_legacy() {
1997 let mut rel = make_rel_table();
1998 rel.insert_row(vec![Value::Int64(2022), Value::Double(3.0)]).unwrap();
2000 assert_eq!(rel.num_rows, 1);
2001 assert_eq!(rel.edges[0], (0, 0)); assert_eq!(rel.get_edge_properties(0), vec![Value::Int64(2022), Value::Double(3.0)]);
2003 }
2004
2005 #[test]
2006 fn test_rel_wrong_column_count() {
2007 let mut rel = make_rel_table();
2008 let result = rel.insert_rel(0, 1, vec![Value::Int64(42)]); assert!(result.is_err());
2010 }
2011
2012 #[test]
2013 fn test_rel_get_column() {
2014 let mut rel = make_rel_table();
2015 rel.insert_rel(0, 1, vec![Value::Int64(2020), Value::Double(1.5)])
2016 .unwrap();
2017 rel.insert_rel(1, 2, vec![Value::Int64(2021), Value::Double(2.5)])
2018 .unwrap();
2019
2020 let since_col = rel.get_column(0).unwrap();
2021 assert_eq!(since_col, &[Value::Int64(2020), Value::Int64(2021)]);
2022
2023 let weight_col = rel.get_column(1).unwrap();
2024 assert_eq!(weight_col, &[Value::Double(1.5), Value::Double(2.5)]);
2025 }
2026
2027 #[test]
2028 fn test_rel_to_column_major() {
2029 let mut rel = make_rel_table();
2030 rel.insert_rel(0, 1, vec![Value::Int64(2020), Value::Double(0.5)])
2031 .unwrap();
2032 rel.insert_rel(2, 3, vec![Value::Int64(2021), Value::Double(0.9)])
2033 .unwrap();
2034
2035 let data = rel.to_column_major_data();
2036 assert_eq!(data.len(), 2);
2037 assert_eq!(data[0], vec![Value::Int64(2020), Value::Int64(2021)]);
2038 assert_eq!(data[1], vec![Value::Double(0.5), Value::Double(0.9)]);
2039 }
2040
2041 #[test]
2044 fn test_catalog_create_and_lookup() {
2045 let cat = TableCatalog::new();
2046 let node_table = cat.create_node_table(
2047 "Person".into(),
2048 vec![ColumnDefinition {
2049 compression: akar_common::enums::CompressionType::Uncompressed,
2050 name: "id".into(),
2051 logical_type: LogicalTypeID::Int64,
2052 is_primary_key: true,
2053 }],
2054 );
2055 assert_eq!(node_table.table_id, 0);
2056
2057 let rel_table = cat.create_rel_table(
2058 "Knows".into(),
2059 0,
2060 1,
2061 vec![ColumnDefinition {
2062 compression: akar_common::enums::CompressionType::Uncompressed,
2063 name: "since".into(),
2064 logical_type: LogicalTypeID::Int64,
2065 is_primary_key: false,
2066 }],
2067 );
2068 assert_eq!(rel_table.table_id, 1);
2069
2070 assert!(cat.get_node_table(0).is_some());
2071 assert!(cat.get_rel_table(1).is_some());
2072 assert_eq!(cat.node_table_num_rows("Person"), 0);
2073 }
2074
2075 #[test]
2076 fn test_refresh_vector_index_after_dml() {
2077 let cat = TableCatalog::new();
2081 cat.create_node_table(
2082 "Item".into(),
2083 vec![
2084 ColumnDefinition {
2085 compression: akar_common::enums::CompressionType::Uncompressed,
2086 name: "id".into(),
2087 logical_type: LogicalTypeID::Int64,
2088 is_primary_key: true,
2089 },
2090 ColumnDefinition {
2091 compression: akar_common::enums::CompressionType::Uncompressed,
2092 name: "embedding".into(),
2093 logical_type: LogicalTypeID::List,
2094 is_primary_key: false,
2095 },
2096 ],
2097 );
2098
2099 let vec3 = |x: f64, y: f64, z: f64| Value::List(vec![Value::Double(x), Value::Double(y), Value::Double(z)]);
2100 {
2101 let mut t = cat.get_node_table_by_name_mut("Item").unwrap();
2102 t.insert_row(vec![Value::Int64(1), vec3(1.0, 0.0, 0.0)]).unwrap();
2103 t.insert_row(vec![Value::Int64(2), vec3(0.0, 1.0, 0.0)]).unwrap();
2104 t.insert_row(vec![Value::Int64(3), vec3(0.0, 0.0, 1.0)]).unwrap();
2105 }
2106
2107 cat.create_vector_index(
2108 "item_vec".into(),
2109 "Item".into(),
2110 "embedding".into(),
2111 DistanceMetric::Cosine,
2112 3,
2113 );
2114
2115 cat.refresh_vector_indexes_for_tables(&[0]);
2117 {
2118 let vi = cat.get_vector_index(1).unwrap();
2119 assert_eq!(vi.hnsw().len(), 3);
2120 assert!(vi.hnsw().get_vector(2).is_some());
2121 }
2122
2123 {
2125 let mut t = cat.get_node_table_by_name_mut("Item").unwrap();
2126 t.insert_row(vec![Value::Int64(4), vec3(1.0, 1.0, 0.0)]).unwrap();
2127 t.insert_row(vec![Value::Int64(5), vec3(0.0, 1.0, 1.0)]).unwrap();
2128 }
2129
2130 cat.refresh_vector_indexes_for_tables(&[0]);
2131 let vi = cat.get_vector_index(1).unwrap();
2132 assert_eq!(vi.hnsw().len(), 5);
2133 assert!(vi.hnsw().get_vector(4).is_some());
2134 }
2135}