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::vector_index::VectorIndexTable;
11use akar_common::error::StorageError;
12use akar_common::types::{LogicalTypeID, Value, pk_value_to_string};
13use akar_vector::hnsw::DistanceMetric;
14use dashmap::DashMap;
15use std::collections::HashMap;
16
17#[derive(Debug, Clone)]
19pub struct ColumnDefinition {
20 pub name: String,
21 pub logical_type: LogicalTypeID,
22 pub is_primary_key: bool,
23 pub compression: akar_common::enums::CompressionType,
24}
25
26#[derive(Debug, Clone)]
37pub struct NodeTable {
38 pub table_id: u64,
39 pub name: String,
40 pub columns: Vec<ColumnDefinition>,
41 pub primary_key_column: usize,
42 pub num_rows: u64,
43 pub node_groups: Vec<NodeGroup>,
46 pub hash_index: HashIndex<String>,
49 pub art_index: Option<ArtPrimaryKeyIndex>,
52 pub persistence_dirty: bool,
55}
56
57pub const NO_PRIMARY_KEY: usize = usize::MAX;
63
64impl NodeTable {
65 pub fn new(table_id: u64, name: String, columns: Vec<ColumnDefinition>) -> Self {
66 let primary_key_column = columns.iter().position(|c| c.is_primary_key).unwrap_or(NO_PRIMARY_KEY);
71 Self {
72 table_id,
73 name,
74 columns,
75 primary_key_column,
76 num_rows: 0,
77 node_groups: Vec::new(),
78 hash_index: HashIndex::new(),
79 art_index: None,
80 persistence_dirty: false,
81 }
82 }
83
84 pub fn add_column(&mut self, column: ColumnDefinition) {
90 if self.columns.iter().any(|c| c.name.eq_ignore_ascii_case(&column.name)) {
91 return;
92 }
93 self.columns.push(column);
94 for group in &mut self.node_groups {
95 let mut chunk = ColumnChunk::new();
96 for _ in 0..group.num_nodes {
97 chunk.append(Value::Null);
98 }
99 group.columns.push(chunk);
100 }
101 self.persistence_dirty = true;
102 }
103
104 pub fn insert_row(&mut self, values: Vec<Value>) -> Result<u64, StorageError> {
119 self.insert_row_with_txn(values, None)
120 }
121
122 pub fn insert_row_with_txn(&mut self, mut values: Vec<Value>, txn_id: Option<u64>) -> Result<u64, StorageError> {
124 if values.len() != self.columns.len() {
125 return Err(StorageError::Page(format!(
126 "Column count mismatch: expected {} values, got {}",
127 self.columns.len(),
128 values.len()
129 )));
130 }
131
132 if self.primary_key_column < self.columns.len() {
134 let pk_value = &values[self.primary_key_column];
135 if matches!(pk_value, Value::Null) {
136 return Err(StorageError::Page(format!(
137 "NULL value not allowed for primary key column '{}' in table '{}'",
138 self.columns[self.primary_key_column].name, self.name
139 )));
140 }
141 }
142
143 coerce_values_to_columns(&mut values, &self.columns)?;
148
149 if self.primary_key_column < self.columns.len() {
151 let pk_value = &values[self.primary_key_column];
152 let pk_key = pk_value_to_string(pk_value);
153 if self.hash_index.lookup(&pk_key).is_some() {
154 return Err(StorageError::Index(format!(
155 "Duplicate primary key value: '{pk_key}' in table '{}'",
156 self.name
157 )));
158 }
159 }
160
161 let num_cols = self.columns.len();
163 if self.node_groups.is_empty() || self.node_groups.last().unwrap().is_full() {
164 let start_offset = self.num_rows;
165 let mut new_group = NodeGroup::new(num_cols, start_offset);
166 if txn_id.is_some() {
168 new_group.enable_version_info();
169 }
170 self.node_groups.push(new_group);
171 }
172
173 let current = self.node_groups.last_mut().unwrap();
174 if txn_id.is_some() {
176 current.enable_version_info();
177 }
178 current.append_row_with_txn(values.clone(), txn_id)?;
179 self.num_rows += 1;
180
181 if self.primary_key_column < self.columns.len() {
183 let pk_value = &values[self.primary_key_column];
184 let pk_key = pk_value_to_string(pk_value);
185 self.hash_index.insert(pk_key, self.num_rows - 1);
186
187 if let Some(ref mut art_idx) = self.art_index
189 && let Some(art_key) = ArtKey::from_value(pk_value)
190 {
191 art_idx.insert(&art_key, self.num_rows - 1);
192 }
193 }
194
195 Ok(self.num_rows - 1)
196 }
197
198 pub fn insert_rows_batch(&mut self, rows: &[Vec<Value>]) -> Result<u64, StorageError> {
202 self.insert_rows_batch_with_txn(rows, None)
203 }
204
205 pub fn insert_rows_batch_with_txn(
207 &mut self,
208 rows: &[Vec<Value>],
209 txn_id: Option<u64>,
210 ) -> Result<u64, StorageError> {
211 if rows.is_empty() {
212 return Ok(0);
213 }
214 let num_cols = self.columns.len();
215
216 for (i, row) in rows.iter().enumerate() {
218 if row.len() != num_cols {
219 return Err(StorageError::Page(format!(
220 "Row {} column count mismatch: expected {} values, got {}",
221 i,
222 num_cols,
223 row.len()
224 )));
225 }
226 if self.primary_key_column < num_cols {
228 let pk_value = &row[self.primary_key_column];
229 if matches!(pk_value, Value::Null) {
230 return Err(StorageError::Page(format!(
231 "NULL value not allowed for primary key column '{}' in table '{}'",
232 self.columns[self.primary_key_column].name, self.name
233 )));
234 }
235 }
236 if self.primary_key_column < num_cols {
238 let pk_key = pk_value_to_string(&row[self.primary_key_column]);
239 if self.hash_index.lookup(&pk_key).is_some() {
240 return Err(StorageError::Index(format!(
241 "Duplicate primary key value: '{pk_key}' in table '{}'",
242 self.name
243 )));
244 }
245 }
246 }
247
248 let start_offset = self.num_rows;
249 let total_new = rows.len();
250
251 let mut coerced_rows = rows.to_vec();
255 for row in &mut coerced_rows {
256 coerce_values_to_columns(row, &self.columns)?;
257 }
258 let rows: &[Vec<Value>] = &coerced_rows;
259
260 if self.node_groups.is_empty() || self.node_groups.last().unwrap().is_full() {
262 let off = if self.node_groups.is_empty() {
263 start_offset
264 } else {
265 self.num_rows
266 };
267 let mut new_group = NodeGroup::new(num_cols, off);
268 if txn_id.is_some() {
269 new_group.enable_version_info();
270 }
271 self.node_groups.push(new_group);
272 }
273
274 let mut inserted = 0usize;
276 while inserted < total_new {
277 let current = self.node_groups.last_mut().unwrap();
278 if txn_id.is_some() {
279 current.enable_version_info();
280 }
281 let rem = current.remaining();
282 let take = (total_new - inserted).min(rem);
283 for row in &rows[inserted..inserted + take] {
284 current.append_row_with_txn(row.clone(), txn_id)?;
285 }
286 self.num_rows += take as u64;
287 inserted += take;
288 if inserted < total_new {
289 let off = self.num_rows;
290 let mut new_group = NodeGroup::new(num_cols, off);
291 if txn_id.is_some() {
292 new_group.enable_version_info();
293 }
294 self.node_groups.push(new_group);
295 }
296 }
297
298 for (i, row) in rows.iter().enumerate() {
300 if self.primary_key_column < num_cols {
301 let pk_key = pk_value_to_string(&row[self.primary_key_column]);
302 self.hash_index.insert(pk_key, start_offset + i as u64);
303 if let Some(ref mut art_idx) = self.art_index
304 && let Some(art_key) = ArtKey::from_value(&row[self.primary_key_column])
305 {
306 art_idx.insert(&art_key, start_offset + i as u64);
307 }
308 }
309 }
310
311 Ok(rows.len() as u64)
312 }
313
314 pub fn lookup_by_pk(&self, pk_value: &Value) -> Option<u64> {
319 let pk_key = pk_value_to_string(pk_value);
320 self.hash_index.lookup(&pk_key)
321 }
322
323 pub fn lookup_by_pk_batch(&self, pk_values: &[Value]) -> Vec<Option<u64>> {
330 pk_values
331 .iter()
332 .map(|pk_value| {
333 let pk_key = pk_value_to_string(pk_value);
334 self.hash_index.lookup(&pk_key)
335 })
336 .collect()
337 }
338
339 pub fn lookup_by_pk_range(
347 &self,
348 lower: Option<&Value>,
349 lower_inclusive: bool,
350 upper: Option<&Value>,
351 upper_inclusive: bool,
352 max_results: u64,
353 ) -> Vec<u64> {
354 match &self.art_index {
355 Some(idx) => {
356 let lower_key = lower.and_then(ArtKey::from_value);
357 let upper_key = upper.and_then(ArtKey::from_value);
358 idx.range_scan(
359 lower_key.as_ref(),
360 lower_inclusive,
361 upper_key.as_ref(),
362 upper_inclusive,
363 max_results,
364 )
365 }
366 None => Vec::new(),
367 }
368 }
369
370 pub fn scan_column(
379 &self,
380 col_idx: usize,
381 start: u64,
382 count: u64,
383 snapshot_ts: Option<u64>,
384 commit_history: &[(u64, u64)],
385 ) -> Vec<Value> {
386 if col_idx >= self.columns.len() || start >= self.num_rows {
387 return Vec::new();
388 }
389 let end = (start + count).min(self.num_rows);
390 let mut result = Vec::with_capacity((end - start) as usize);
391
392 let group_start = self.find_group(start);
394 let mut remaining = end - start;
395
396 for g_idx in group_start..self.node_groups.len() {
397 if remaining == 0 {
398 break;
399 }
400 let group = &self.node_groups[g_idx];
401 let local_start = if g_idx == group_start {
402 (start - group.start_offset) as usize
403 } else {
404 0
405 };
406 let available = (group.num_nodes as usize).saturating_sub(local_start);
407 let take = available.min(remaining as usize);
408
409 for row in local_start..local_start + take {
410 let val = group.get_value_with_snapshot(row, col_idx, snapshot_ts, commit_history);
411 match val {
412 Some(v) => result.push(v.clone()),
413 None => result.push(Value::Null),
414 }
415 }
416 remaining -= take as u64;
417 }
418
419 result
420 }
421
422 pub fn load_persisted_rows(&mut self, rows: Vec<Vec<Value>>) -> Result<(), StorageError> {
430 let num_cols = self.columns.len();
431 self.node_groups.clear();
432 self.hash_index.clear();
433
434 let mut offset = 0u64;
435 let mut group = NodeGroup::new(num_cols, offset);
436 for row in &rows {
437 if row.len() != num_cols {
438 return Err(StorageError::Page(format!(
439 "load_persisted_rows: expected {num_cols} values, got {}",
440 row.len()
441 )));
442 }
443 group.append_row_with_txn(row.clone(), None)?;
444 offset += 1;
445 if group.is_full() {
446 self.node_groups.push(group);
447 group = NodeGroup::new(num_cols, offset);
448 }
449 }
450 if group.num_nodes > 0 {
451 self.node_groups.push(group);
452 }
453 self.num_rows = offset;
454
455 if self.primary_key_column < num_cols {
457 for (row_idx, values) in rows.iter().enumerate() {
458 let pk = &values[self.primary_key_column];
459 if matches!(pk, Value::Null) {
460 continue;
461 }
462 let key = pk_value_to_string(pk);
463 self.hash_index.insert(key, row_idx as u64);
464 }
465 }
466
467 if self.art_index.is_some() && self.primary_key_column < num_cols {
469 if let Some(art) = &mut self.art_index {
470 art.clear();
471 for (row_idx, values) in rows.iter().enumerate() {
472 let pk = &values[self.primary_key_column];
473 if matches!(pk, Value::Null) {
474 continue;
475 }
476 if let Some(key) = ArtKey::from_value(pk) {
477 art.insert(&key, row_idx as u64);
478 }
479 }
480 }
481 }
482 Ok(())
483 }
484
485 pub fn update_cell(&mut self, row_idx: u64, col_idx: usize, value: Value) -> Result<(), StorageError> {
487 if col_idx >= self.columns.len() {
488 return Err(StorageError::Page(format!("Column index {col_idx} out of range")));
489 }
490 if row_idx >= self.num_rows {
491 return Err(StorageError::Page(format!(
492 "Row index {row_idx} out of range (num_rows={})",
493 self.num_rows
494 )));
495 }
496 self.persistence_dirty = true;
497
498 let mut offset = 0u64;
499 for group in &mut self.node_groups {
500 if row_idx < offset + group.num_nodes {
501 let local_row = (row_idx - offset) as usize;
502 if let Some(col_chunk) = group.columns.get_mut(col_idx) {
503 col_chunk.set_value(local_row, value)?;
504 }
505 return Ok(());
506 }
507 offset += group.num_nodes;
508 }
509 Err(StorageError::Page(format!(
510 "Row index {row_idx} not found in any node group"
511 )))
512 }
513
514 pub fn delete_row(&mut self, row_idx: u64) -> Result<(), StorageError> {
517 self.delete_row_with_txn(row_idx, None)
518 }
519
520 pub fn delete_row_with_txn(&mut self, row_idx: u64, txn_id: Option<u64>) -> Result<(), StorageError> {
525 if row_idx >= self.num_rows {
526 return Err(StorageError::Page(format!(
527 "Row index {row_idx} out of range (num_rows={})",
528 self.num_rows
529 )));
530 }
531 self.persistence_dirty = true;
532
533 let mut offset = 0u64;
535 for group in &mut self.node_groups {
536 if row_idx < offset + group.num_nodes {
537 let local_row = (row_idx - offset) as usize;
538 if let Some(txn) = txn_id {
540 if let Some(ref vi) = group.version_info {
541 vi.delete(txn, local_row as u32);
542 }
543 }
544 let pk_value = (self.primary_key_column < self.columns.len())
547 .then(|| group.columns.get(self.primary_key_column))
548 .flatten()
549 .and_then(|chunk| chunk.get(local_row))
550 .cloned();
551 for col_chunk in &mut group.columns {
553 let _ = col_chunk.set_value(local_row, Value::Null);
554 }
555 if let Some(pk) = pk_value
559 && !matches!(pk, Value::Null)
560 {
561 let pk_key = pk_value_to_string(&pk);
562 self.hash_index.delete(&pk_key);
563 if let Some(ref mut art_idx) = self.art_index
564 && let Some(art_key) = ArtKey::from_value(&pk)
565 {
566 art_idx.delete(&art_key, row_idx);
567 }
568 }
569 return Ok(());
570 }
571 offset += group.num_nodes;
572 }
573 Err(StorageError::Page(format!(
574 "Row index {row_idx} not found in any node group"
575 )))
576 }
577
578 pub fn row_undo_bytes(&self, row_idx: u64) -> Vec<u8> {
582 let mut out = Vec::new();
583 for col in 0..self.columns.len() {
584 let val = self.get_value(row_idx as usize, col).cloned().unwrap_or(Value::Null);
585 out.extend_from_slice(&Column::serialize_value(&val));
586 }
587 out
588 }
589
590 pub fn cell_undo_bytes(&self, row_idx: u64, col_idx: usize) -> Vec<u8> {
593 let val = self
594 .get_value(row_idx as usize, col_idx)
595 .cloned()
596 .unwrap_or(Value::Null);
597 Column::serialize_value(&val)
598 }
599
600 pub fn get_value(&self, row: usize, col: usize) -> Option<&Value> {
603 self.get_value_with_snapshot(row, col, None, &[])
604 }
605
606 pub fn get_value_with_snapshot(
611 &self,
612 row: usize,
613 col: usize,
614 snapshot_ts: Option<u64>,
615 commit_history: &[(u64, u64)],
616 ) -> Option<&Value> {
617 if col >= self.columns.len() || row as u64 >= self.num_rows {
618 return None;
619 }
620 let group_idx = self.find_group(row as u64);
621 let group = self.node_groups.get(group_idx)?;
622 let local_row = row as u64 - group.start_offset;
623 group.get_value_with_snapshot(local_row as usize, col, snapshot_ts, commit_history)
624 }
625
626 pub fn to_column_major_data(&self) -> Vec<Vec<Value>> {
630 self.to_column_major_data_with_predicate(None)
631 }
632
633 pub fn to_column_major_data_with_predicate(&self, predicate: Option<(usize, &str, &Value)>) -> Vec<Vec<Value>> {
636 let num_cols = self.columns.len();
637 let mut result = vec![Vec::new(); num_cols]; for group in &self.node_groups {
640 if let Some((col_idx, op, val)) = predicate
641 && let Some(col_chunk) = group.columns.get(col_idx)
642 {
643 use crate::predicate::{ZoneMapCheckResult, check_zone_map};
644 if check_zone_map(&col_chunk.stats, op, val) == ZoneMapCheckResult::SkipScan {
645 continue; }
647 }
648
649 for row in 0..group.num_nodes as usize {
650 for (col, res_col) in result.iter_mut().enumerate().take(num_cols) {
651 match group.get_value(row, col) {
652 Some(v) => res_col.push(v.clone()),
653 None => res_col.push(Value::Null),
654 }
655 }
656 }
657 }
658
659 result
660 }
661
662 pub fn to_column_major_data_with_snapshot(
668 &self,
669 snapshot_ts: Option<u64>,
670 commit_history: &[(u64, u64)],
671 ) -> Vec<Vec<Value>> {
672 let num_cols = self.columns.len();
673 let mut result = vec![Vec::new(); num_cols];
674
675 for group in &self.node_groups {
676 for row in 0..group.num_nodes as usize {
677 for (col, res_col) in result.iter_mut().enumerate().take(num_cols) {
678 match group.get_value_owned_with_snapshot(row, col, snapshot_ts, commit_history) {
679 Some(v) => res_col.push(v),
680 None => res_col.push(Value::Null),
681 }
682 }
683 }
684 }
685
686 result
687 }
688
689 pub fn to_column_major_data_with_snapshot_and_predicate(
692 &self,
693 predicate: Option<(usize, &str, &Value)>,
694 snapshot_ts: Option<u64>,
695 commit_history: &[(u64, u64)],
696 ) -> Vec<Vec<Value>> {
697 let num_cols = self.columns.len();
698 let mut result = vec![Vec::new(); num_cols];
699
700 for group in &self.node_groups {
701 if let Some((col_idx, op, val)) = predicate
702 && let Some(col_chunk) = group.columns.get(col_idx)
703 {
704 use crate::predicate::{ZoneMapCheckResult, check_zone_map};
705 if check_zone_map(&col_chunk.stats, op, val) == ZoneMapCheckResult::SkipScan {
706 continue;
707 }
708 }
709
710 for row in 0..group.num_nodes as usize {
711 for (col, res_col) in result.iter_mut().enumerate().take(num_cols) {
712 match group.get_value_owned_with_snapshot(row, col, snapshot_ts, commit_history) {
713 Some(v) => res_col.push(v),
714 None => res_col.push(Value::Null),
715 }
716 }
717 }
718 }
719
720 result
721 }
722
723 pub fn to_column_major_data_with_predicate_and_ids(
728 &self,
729 predicate: Option<(usize, &str, &Value)>,
730 ) -> (Vec<Vec<Value>>, Vec<u64>) {
731 let num_cols = self.columns.len();
732 let mut result = vec![Vec::new(); num_cols];
733 let mut ids = Vec::new();
734
735 for group in &self.node_groups {
736 if let Some((col_idx, op, val)) = predicate
737 && let Some(col_chunk) = group.columns.get(col_idx)
738 {
739 use crate::predicate::{ZoneMapCheckResult, check_zone_map};
740 if check_zone_map(&col_chunk.stats, op, val) == ZoneMapCheckResult::SkipScan {
741 continue;
742 }
743 }
744
745 for row in 0..group.num_nodes as usize {
746 ids.push(group.start_offset + row as u64);
747 for (col, res_col) in result.iter_mut().enumerate().take(num_cols) {
748 match group.get_value(row, col) {
749 Some(v) => res_col.push(v.clone()),
750 None => res_col.push(Value::Null),
751 }
752 }
753 }
754 }
755
756 (result, ids)
757 }
758
759 pub fn to_column_major_data_with_snapshot_and_predicate_and_ids(
763 &self,
764 predicate: Option<(usize, &str, &Value)>,
765 snapshot_ts: Option<u64>,
766 commit_history: &[(u64, u64)],
767 ) -> (Vec<Vec<Value>>, Vec<u64>) {
768 let num_cols = self.columns.len();
769 let mut result = vec![Vec::new(); num_cols];
770 let mut ids = Vec::new();
771
772 for group in &self.node_groups {
773 if let Some((col_idx, op, val)) = predicate
774 && let Some(col_chunk) = group.columns.get(col_idx)
775 {
776 use crate::predicate::{ZoneMapCheckResult, check_zone_map};
777 if check_zone_map(&col_chunk.stats, op, val) == ZoneMapCheckResult::SkipScan {
778 continue;
779 }
780 }
781
782 for row in 0..group.num_nodes as usize {
783 ids.push(group.start_offset + row as u64);
784 for (col, res_col) in result.iter_mut().enumerate().take(num_cols) {
785 match group.get_value_owned_with_snapshot(row, col, snapshot_ts, commit_history) {
786 Some(v) => res_col.push(v),
787 None => res_col.push(Value::Null),
788 }
789 }
790 }
791 }
792
793 (result, ids)
794 }
795
796 fn find_group(&self, row: u64) -> usize {
798 self.node_groups
799 .binary_search_by_key(&row, |g| g.start_offset)
800 .unwrap_or_else(|i| if i == 0 { 0 } else { i - 1 })
801 }
802}
803
804fn coerce_values_to_columns(values: &mut [Value], columns: &[ColumnDefinition]) -> Result<(), StorageError> {
812 for (i, col) in columns.iter().enumerate() {
813 let coerced = match (col.logical_type, &values[i]) {
814 (LogicalTypeID::UInt64, Value::Int64(x)) if *x >= 0 => Some(Value::UInt64(*x as u64)),
815 (LogicalTypeID::UInt64, Value::Int32(x)) if *x >= 0 => Some(Value::UInt64(*x as u64)),
816 (LogicalTypeID::UInt64, Value::Int16(x)) if *x >= 0 => Some(Value::UInt64(*x as u64)),
817 (LogicalTypeID::UInt64, Value::Int8(x)) if *x >= 0 => Some(Value::UInt64(*x as u64)),
818 (LogicalTypeID::UInt64, Value::UInt32(x)) => Some(Value::UInt64(*x as u64)),
819 (LogicalTypeID::UInt64, Value::UInt16(x)) => Some(Value::UInt64(*x as u64)),
820 (LogicalTypeID::UInt64, Value::UInt8(x)) => Some(Value::UInt64(*x as u64)),
821 (LogicalTypeID::UInt64, Value::Int64(_))
822 | (LogicalTypeID::UInt64, Value::Int32(_))
823 | (LogicalTypeID::UInt64, Value::Int16(_))
824 | (LogicalTypeID::UInt64, Value::Int8(_)) => {
825 return Err(StorageError::Page(format!(
826 "Cannot store negative value in UINT64 column '{}'",
827 col.name
828 )));
829 }
830 _ => None,
831 };
832 if let Some(c) = coerced {
833 values[i] = c;
834 }
835 }
836 Ok(())
837}
838
839#[derive(Debug, Clone)]
851pub struct RelTable {
852 pub table_id: u64,
853 pub name: String,
854 pub src_table_id: u64,
855 pub dst_table_id: u64,
856 pub columns: Vec<ColumnDefinition>,
857 pub num_rows: u64,
858 pub edges: Vec<(u64, u64)>,
860 pub fwd_adj: HashMap<u64, Vec<(u64, usize)>>,
862 pub rev_adj: HashMap<u64, Vec<(u64, usize)>>,
864 pub csr_index: Option<CsrIndex>,
866 pub properties: Vec<Vec<Value>>,
868 pub persistence_dirty: bool,
871}
872
873impl RelTable {
874 pub fn new(
875 table_id: u64,
876 name: String,
877 src_table_id: u64,
878 dst_table_id: u64,
879 columns: Vec<ColumnDefinition>,
880 ) -> Self {
881 let num_cols = columns.len();
882 Self {
883 table_id,
884 name,
885 src_table_id,
886 dst_table_id,
887 columns,
888 num_rows: 0,
889 edges: Vec::new(),
890 fwd_adj: HashMap::new(),
891 rev_adj: HashMap::new(),
892 csr_index: None,
893 properties: vec![Vec::new(); num_cols],
894 persistence_dirty: false,
895 }
896 }
897
898 pub fn add_column(&mut self, column: ColumnDefinition) {
901 if self.columns.iter().any(|c| c.name.eq_ignore_ascii_case(&column.name)) {
902 return;
903 }
904 self.columns.push(column);
905 self.properties.push(vec![Value::Null; self.edges.len()]);
906 self.persistence_dirty = true;
907 }
908
909 pub fn insert_rel(&mut self, from: u64, to: u64, values: Vec<Value>) -> Result<(), StorageError> {
917 if values.len() != self.columns.len() {
918 return Err(StorageError::Page(format!(
919 "Column count mismatch: expected {} values, got {}",
920 self.columns.len(),
921 values.len()
922 )));
923 }
924
925 let edge_idx = self.edges.len();
926 self.edges.push((from, to));
927
928 self.fwd_adj.entry(from).or_default().push((to, edge_idx));
930
931 self.rev_adj.entry(to).or_default().push((from, edge_idx));
933 for (col_idx, val) in values.into_iter().enumerate() {
935 self.properties[col_idx].push(val);
936 }
937 self.num_rows += 1;
938 Ok(())
939 }
940
941 pub fn insert_rels_batch(&mut self, rels: &[(u64, u64, Vec<Value>)]) -> Result<u64, StorageError> {
944 if rels.is_empty() {
945 return Ok(0);
946 }
947 let num_cols = self.columns.len();
948 let total = rels.len();
949
950 for (i, (_, _, vals)) in rels.iter().enumerate() {
952 if vals.len() != num_cols {
953 return Err(StorageError::Page(format!(
954 "Rel {} column count mismatch: expected {} values, got {}",
955 i,
956 num_cols,
957 vals.len()
958 )));
959 }
960 }
961
962 self.edges.reserve(total);
964 for col in &mut self.properties {
965 col.reserve(total);
966 }
967
968 let _start_edge_idx = self.edges.len();
969
970 for (from, to, vals) in rels {
972 let edge_idx = self.edges.len();
973 self.edges.push((*from, *to));
974 self.fwd_adj.entry(*from).or_default().push((*to, edge_idx));
975 self.rev_adj.entry(*to).or_default().push((*from, edge_idx));
976 for (col_idx, val) in vals.iter().enumerate() {
977 self.properties[col_idx].push(val.clone());
978 }
979 }
980
981 self.num_rows += total as u64;
982 Ok(total as u64)
983 }
984
985 pub fn delete_edge(&mut self, edge_idx: usize) -> Result<(), StorageError> {
988 if edge_idx >= self.edges.len() {
989 return Err(StorageError::Page(format!("Edge index {edge_idx} out of range")));
990 }
991
992 let (src, dst) = self.edges[edge_idx];
993 if src == u64::MAX {
994 return Ok(());
996 }
997
998 if let Some(adj) = self.fwd_adj.get_mut(&src) {
1000 adj.retain(|&(_, idx)| idx != edge_idx);
1001 }
1002
1003 if let Some(adj) = self.rev_adj.get_mut(&dst) {
1005 adj.retain(|&(_, idx)| idx != edge_idx);
1006 }
1007
1008 self.edges[edge_idx] = (u64::MAX, u64::MAX);
1010 self.persistence_dirty = true;
1011
1012 for col in &mut self.properties {
1014 if edge_idx < col.len() {
1015 col[edge_idx] = Value::Null;
1016 }
1017 }
1018
1019 Ok(())
1020 }
1021
1022 pub fn update_cell(&mut self, edge_idx: usize, col_idx: usize, value: Value) -> Result<(), StorageError> {
1024 if col_idx >= self.columns.len() {
1025 return Err(StorageError::Page(format!("Column index {col_idx} out of range")));
1026 }
1027 if edge_idx >= self.properties[col_idx].len() {
1028 return Err(StorageError::Page(format!("Edge index {edge_idx} out of range")));
1029 }
1030
1031 self.properties[col_idx][edge_idx] = value;
1032 self.persistence_dirty = true;
1033 Ok(())
1034 }
1035
1036 pub fn edge_undo_bytes(&self, edge_idx: usize) -> Vec<u8> {
1040 let (src, dst) = self.edges.get(edge_idx).copied().unwrap_or((u64::MAX, u64::MAX));
1041 let mut out = Vec::new();
1042 out.extend_from_slice(&Column::serialize_value(&Value::UInt64(src)));
1043 out.extend_from_slice(&Column::serialize_value(&Value::UInt64(dst)));
1044 for p in self.get_edge_properties(edge_idx) {
1045 out.extend_from_slice(&Column::serialize_value(&p));
1046 }
1047 out
1048 }
1049
1050 pub fn edge_cell_undo_bytes(&self, edge_idx: usize, col_idx: usize) -> Vec<u8> {
1053 let v = self
1054 .get_edge_properties(edge_idx)
1055 .get(col_idx)
1056 .cloned()
1057 .unwrap_or(Value::Null);
1058 Column::serialize_value(&v)
1059 }
1060
1061 pub fn restore_deleted_edge(
1065 &mut self,
1066 edge_idx: usize,
1067 src: u64,
1068 dst: u64,
1069 props: Vec<Value>,
1070 ) -> Result<(), StorageError> {
1071 if edge_idx >= self.edges.len() {
1072 return Err(StorageError::Page(format!("Edge index {edge_idx} out of range")));
1073 }
1074 self.edges[edge_idx] = (src, dst);
1075 self.fwd_adj.entry(src).or_default().push((dst, edge_idx));
1076 self.rev_adj.entry(dst).or_default().push((src, edge_idx));
1077 for (col_idx, val) in props.into_iter().enumerate() {
1078 if col_idx < self.properties.len() {
1079 if edge_idx < self.properties[col_idx].len() {
1080 self.properties[col_idx][edge_idx] = val;
1081 } else {
1082 self.properties[col_idx].push(val);
1083 }
1084 }
1085 }
1086 self.persistence_dirty = true;
1087 Ok(())
1088 }
1089
1090 pub fn insert_row(&mut self, values: Vec<Value>) -> Result<u64, StorageError> {
1094 let num_prop_cols = self.columns.len();
1098 if values.len() != num_prop_cols {
1099 return Err(StorageError::Page(format!(
1100 "Column count mismatch: expected {} values, got {}",
1101 num_prop_cols,
1102 values.len()
1103 )));
1104 }
1105
1106 let from = self.num_rows;
1109 let to = self.num_rows;
1110 self.insert_rel(from, to, values)?;
1111 Ok(0) }
1113
1114 pub fn scan_adj_list(&self, src_offset: u64) -> &[(u64, usize)] {
1119 self.fwd_adj.get(&src_offset).map(|v| v.as_slice()).unwrap_or(&[])
1120 }
1121
1122 pub fn scan_rev_adj_list(&self, dst_offset: u64) -> &[(u64, usize)] {
1127 self.rev_adj.get(&dst_offset).map(|v| v.as_slice()).unwrap_or(&[])
1128 }
1129
1130 pub fn get_outgoing_edges(&self, src_offset: u64) -> Vec<(u64, Vec<Value>)> {
1132 self.scan_adj_list(src_offset)
1133 .iter()
1134 .map(|&(dst, edge_idx)| {
1135 let props = self.get_edge_properties(edge_idx);
1136 (dst, props)
1137 })
1138 .collect()
1139 }
1140
1141 pub fn get_incoming_edges(&self, dst_offset: u64) -> Vec<(u64, Vec<Value>)> {
1143 self.scan_rev_adj_list(dst_offset)
1144 .iter()
1145 .map(|&(src, edge_idx)| {
1146 let props = self.get_edge_properties(edge_idx);
1147 (src, props)
1148 })
1149 .collect()
1150 }
1151
1152 pub fn get_edge_properties(&self, edge_idx: usize) -> Vec<Value> {
1154 let mut props = Vec::with_capacity(self.columns.len());
1155 for col in &self.properties {
1156 match col.get(edge_idx) {
1157 Some(v) => props.push(v.clone()),
1158 None => props.push(Value::Null),
1159 }
1160 }
1161 props
1162 }
1163
1164 pub fn get_column(&self, col_idx: usize) -> Option<&[Value]> {
1166 self.properties.get(col_idx).map(|v| v.as_slice())
1167 }
1168
1169 pub fn to_column_major_data(&self) -> Vec<Vec<Value>> {
1171 self.properties.clone()
1172 }
1173}
1174
1175#[derive(Debug, Default)]
1182pub struct TableCatalog {
1183 node_tables: DashMap<u64, NodeTable>,
1184 rel_tables: DashMap<u64, RelTable>,
1185 vector_indexes: DashMap<u64, VectorIndexTable>,
1186 node_name_to_id: DashMap<String, u64>,
1188 rel_name_to_id: DashMap<String, u64>,
1190 vector_index_name_to_id: DashMap<String, u64>,
1192 next_table_id: std::sync::atomic::AtomicU64,
1193}
1194
1195impl TableCatalog {
1196 pub fn new() -> Self {
1197 Self::default()
1198 }
1199
1200 pub fn create_node_table(&self, name: String, columns: Vec<ColumnDefinition>) -> NodeTable {
1201 let table_id = self.next_table_id.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1202 let table = NodeTable::new(table_id, name.clone(), columns);
1203 self.node_name_to_id.insert(name, table_id);
1204 self.node_tables.insert(table_id, table.clone());
1205 table
1206 }
1207
1208 pub fn create_node_table_with_id(&self, table_id: u64, name: String, columns: Vec<ColumnDefinition>) -> NodeTable {
1212 self.bump_next_table_id(table_id);
1213 let table = NodeTable::new(table_id, name.clone(), columns);
1214 self.node_name_to_id.insert(name, table_id);
1215 self.node_tables.insert(table_id, table.clone());
1216 table
1217 }
1218
1219 pub fn create_rel_table(
1220 &self,
1221 name: String,
1222 src_table_id: u64,
1223 dst_table_id: u64,
1224 columns: Vec<ColumnDefinition>,
1225 ) -> RelTable {
1226 let table_id = self.next_table_id.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1227 let table = RelTable::new(table_id, name.clone(), src_table_id, dst_table_id, columns);
1228 self.rel_name_to_id.insert(name, table_id);
1229 self.rel_tables.insert(table_id, table.clone());
1230 table
1231 }
1232
1233 pub fn create_rel_table_with_id(
1237 &self,
1238 table_id: u64,
1239 name: String,
1240 src_table_id: u64,
1241 dst_table_id: u64,
1242 columns: Vec<ColumnDefinition>,
1243 ) -> RelTable {
1244 self.bump_next_table_id(table_id);
1245 let table = RelTable::new(table_id, name.clone(), src_table_id, dst_table_id, columns);
1246 self.rel_name_to_id.insert(name, table_id);
1247 self.rel_tables.insert(table_id, table.clone());
1248 table
1249 }
1250
1251 fn bump_next_table_id(&self, table_id: u64) {
1254 let mut next = self.next_table_id.load(std::sync::atomic::Ordering::SeqCst);
1255 while next <= table_id {
1256 match self.next_table_id.compare_exchange(
1257 next,
1258 table_id + 1,
1259 std::sync::atomic::Ordering::SeqCst,
1260 std::sync::atomic::Ordering::SeqCst,
1261 ) {
1262 Ok(_) => break,
1263 Err(current) => next = current,
1264 }
1265 }
1266 }
1267
1268 pub fn get_node_table(&self, table_id: u64) -> Option<dashmap::mapref::one::Ref<'_, u64, NodeTable>> {
1269 self.node_tables.get(&table_id)
1270 }
1271
1272 pub fn get_node_table_mut(&self, table_id: u64) -> Option<dashmap::mapref::one::RefMut<'_, u64, NodeTable>> {
1273 self.node_tables.get_mut(&table_id)
1274 }
1275
1276 pub fn get_node_table_by_name(&self, name: &str) -> Option<dashmap::mapref::one::Ref<'_, u64, NodeTable>> {
1277 let id = self.node_name_to_id.get(name)?;
1278 self.node_tables.get(&*id)
1279 }
1280
1281 pub fn get_node_table_by_name_mut(&self, name: &str) -> Option<dashmap::mapref::one::RefMut<'_, u64, NodeTable>> {
1282 let id = self.node_name_to_id.get(name)?;
1283 self.node_tables.get_mut(&*id)
1284 }
1285
1286 pub fn get_rel_table(&self, table_id: u64) -> Option<dashmap::mapref::one::Ref<'_, u64, RelTable>> {
1287 self.rel_tables.get(&table_id)
1288 }
1289
1290 pub fn get_rel_table_mut(&self, table_id: u64) -> Option<dashmap::mapref::one::RefMut<'_, u64, RelTable>> {
1291 self.rel_tables.get_mut(&table_id)
1292 }
1293
1294 pub fn get_rel_table_by_name(&self, name: &str) -> Option<dashmap::mapref::one::Ref<'_, u64, RelTable>> {
1295 let id = self.rel_name_to_id.get(name)?;
1296 self.rel_tables.get(&*id)
1297 }
1298
1299 pub fn get_rel_table_by_name_mut(&self, name: &str) -> Option<dashmap::mapref::one::RefMut<'_, u64, RelTable>> {
1300 let id = self.rel_name_to_id.get(name)?;
1301 self.rel_tables.get_mut(&*id)
1302 }
1303
1304 pub fn has_incident_edges(&self, table_id: u64, node_idx: u64) -> bool {
1306 for rel_table in self.rel_tables.iter() {
1307 if rel_table.src_table_id == table_id {
1308 if let Some(edges) = rel_table.fwd_adj.get(&node_idx) {
1309 if !edges.is_empty() {
1310 return true;
1311 }
1312 }
1313 }
1314 if rel_table.dst_table_id == table_id {
1315 if let Some(edges) = rel_table.rev_adj.get(&node_idx) {
1316 if !edges.is_empty() {
1317 return true;
1318 }
1319 }
1320 }
1321 }
1322 false
1323 }
1324
1325 pub fn detach_node(&self, table_id: u64, node_idx: u64) {
1327 for mut rel_table in self.rel_tables.iter_mut() {
1328 let mut edges_to_delete = Vec::new();
1329
1330 if rel_table.src_table_id == table_id {
1331 if let Some(edges) = rel_table.fwd_adj.get(&node_idx) {
1332 for &(_, edge_idx) in edges {
1333 edges_to_delete.push(edge_idx);
1334 }
1335 }
1336 }
1337
1338 if rel_table.dst_table_id == table_id {
1339 if let Some(edges) = rel_table.rev_adj.get(&node_idx) {
1340 for &(_, edge_idx) in edges {
1341 edges_to_delete.push(edge_idx);
1342 }
1343 }
1344 }
1345
1346 for edge_idx in edges_to_delete {
1347 let _ = rel_table.delete_edge(edge_idx);
1348 }
1349 }
1350 }
1351
1352 pub fn all_node_tables(&self) -> Vec<dashmap::mapref::multiple::RefMulti<'_, u64, NodeTable>> {
1353 self.node_tables.iter().collect()
1354 }
1355
1356 pub fn all_rel_tables(&self) -> Vec<dashmap::mapref::multiple::RefMulti<'_, u64, RelTable>> {
1357 self.rel_tables.iter().collect()
1358 }
1359
1360 pub fn node_table_num_rows(&self, name: &str) -> u64 {
1362 self.get_node_table_by_name(name).map(|t| t.num_rows).unwrap_or(0)
1363 }
1364
1365 pub fn drop_node_table(&self, name: &str) -> bool {
1367 if let Some(id) = self.node_name_to_id.get(name) {
1368 let table_id = *id;
1369 drop(id);
1370 self.node_name_to_id.remove(name);
1371 self.node_tables.remove(&table_id).is_some()
1372 } else {
1373 false
1374 }
1375 }
1376
1377 pub fn create_vector_index(
1379 &self,
1380 name: String,
1381 table_name: String,
1382 column_name: String,
1383 metric: DistanceMetric,
1384 dimensions: u32,
1385 ) -> VectorIndexTable {
1386 let index_id = self.next_table_id.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1387 let table = VectorIndexTable::new(index_id, name.clone(), table_name, column_name, metric, dimensions);
1388 self.vector_index_name_to_id.insert(name, index_id);
1389 self.vector_indexes.insert(index_id, table.clone());
1390 table
1391 }
1392
1393 pub fn get_vector_index(&self, index_id: u64) -> Option<dashmap::mapref::one::Ref<'_, u64, VectorIndexTable>> {
1395 self.vector_indexes.get(&index_id)
1396 }
1397
1398 pub fn get_vector_index_by_name(&self, name: &str) -> Option<dashmap::mapref::one::Ref<'_, u64, VectorIndexTable>> {
1400 let id = self.vector_index_name_to_id.get(name)?;
1401 self.vector_indexes.get(&*id)
1402 }
1403
1404 pub fn get_vector_index_by_name_mut(
1406 &self,
1407 name: &str,
1408 ) -> Option<dashmap::mapref::one::RefMut<'_, u64, VectorIndexTable>> {
1409 let id = self.vector_index_name_to_id.get(name)?;
1410 self.vector_indexes.get_mut(&*id)
1411 }
1412
1413 pub fn get_vector_index_mut(
1415 &self,
1416 index_id: u64,
1417 ) -> Option<dashmap::mapref::one::RefMut<'_, u64, VectorIndexTable>> {
1418 self.vector_indexes.get_mut(&index_id)
1419 }
1420
1421 pub fn drop_vector_index(&self, name: &str) -> bool {
1423 if let Some(id) = self.vector_index_name_to_id.get(name) {
1424 let index_id = *id;
1425 drop(id);
1426 self.vector_index_name_to_id.remove(name);
1427 self.vector_indexes.remove(&index_id).is_some()
1428 } else {
1429 false
1430 }
1431 }
1432
1433 pub fn all_vector_indexes(&self) -> Vec<dashmap::mapref::multiple::RefMulti<'_, u64, VectorIndexTable>> {
1435 self.vector_indexes.iter().collect()
1436 }
1437
1438 pub fn refresh_vector_index(&self, index_id: u64) {
1443 let (table_name, column_name) = match self.vector_indexes.get(&index_id) {
1444 Some(vi) => (vi.table_name.clone(), vi.column_name.clone()),
1445 None => return,
1446 };
1447
1448 let col_idx = match self.get_node_table_by_name(&table_name) {
1449 Some(t) => match t.columns.iter().position(|c| c.name == column_name) {
1450 Some(idx) => idx,
1451 None => return,
1452 },
1453 None => return,
1454 };
1455
1456 let mut data: Vec<(usize, Vec<f64>)> = Vec::new();
1459 if let Some(table) = self.get_node_table_by_name(&table_name) {
1460 for row_id in 0..table.num_rows as usize {
1461 if let Some(val) = table.get_value(row_id, col_idx) {
1462 if let Ok(vec) = crate::extract_f64_list_from_value(val) {
1463 data.push((row_id, vec));
1464 }
1465 }
1466 }
1467 }
1468 if data.is_empty() {
1469 if let Some(mut vi) = self.vector_indexes.get_mut(&index_id) {
1470 vi.hnsw_mut().clear();
1471 }
1472 return;
1473 }
1474
1475 let mut vi = self.vector_indexes.get_mut(&index_id);
1476 if let Some(vi) = vi.as_mut() {
1477 vi.hnsw_mut().clear();
1478 for (row_id, vec) in data {
1479 vi.hnsw_mut().insert(vec, row_id);
1480 }
1481 }
1482 }
1483
1484 pub fn refresh_vector_indexes_for_tables(&self, table_ids: &[u64]) {
1489 for table_id in table_ids {
1490 let table_name = match self.get_node_table(*table_id) {
1491 Some(t) => t.name.clone(),
1492 None => continue,
1493 };
1494 let index_ids: Vec<u64> = self
1495 .vector_indexes
1496 .iter()
1497 .filter(|vi| vi.table_name == table_name)
1498 .map(|vi| *vi.key())
1499 .collect();
1500 for index_id in index_ids {
1501 self.refresh_vector_index(index_id);
1502 }
1503 }
1504 }
1505
1506 pub fn create_art_index(&self, table_name: &str, index_name: &str) -> Result<(), StorageError> {
1513 let mut table = self
1514 .get_node_table_by_name_mut(table_name)
1515 .ok_or_else(|| StorageError::TableNotFound(format!("Node table '{table_name}' not found")))?;
1516
1517 if table.art_index.is_some() {
1518 return Err(StorageError::Index(format!(
1519 "Table '{table_name}' already has an ART index"
1520 )));
1521 }
1522
1523 let mut art_idx = ArtPrimaryKeyIndex::new(index_name);
1524
1525 let pk_col = table.primary_key_column;
1527 let col_major = table.to_column_major_data();
1529 if pk_col < col_major.len() {
1530 for (row_offset, pk_val) in col_major[pk_col].iter().enumerate() {
1531 if !matches!(pk_val, Value::Null)
1532 && let Some(art_key) = ArtKey::from_value(pk_val)
1533 {
1534 art_idx.insert(&art_key, row_offset as u64);
1535 }
1536 }
1537 }
1538
1539 table.art_index = Some(art_idx);
1540 Ok(())
1541 }
1542
1543 pub fn drop_art_index(&self, table_name: &str) -> Result<(), StorageError> {
1545 let mut table = self
1546 .get_node_table_by_name_mut(table_name)
1547 .ok_or_else(|| StorageError::TableNotFound(format!("Node table '{table_name}' not found")))?;
1548
1549 table.art_index = None;
1550 Ok(())
1551 }
1552
1553 pub fn get_art_index(&self, table_name: &str) -> Option<ArtPrimaryKeyIndex> {
1556 let table = self.get_node_table_by_name(table_name)?;
1557 table.art_index.clone()
1558 }
1559
1560 pub fn has_art_index(&self, table_name: &str) -> bool {
1562 self.get_node_table_by_name(table_name)
1563 .map(|t| t.art_index.is_some())
1564 .unwrap_or(false)
1565 }
1566
1567 pub fn drop_rel_table(&self, name: &str) -> bool {
1569 if let Some(id) = self.rel_name_to_id.get(name) {
1570 let table_id = *id;
1571 drop(id);
1572 self.rel_name_to_id.remove(name);
1573 self.rel_tables.remove(&table_id).is_some()
1574 } else {
1575 false
1576 }
1577 }
1578}
1579
1580#[cfg(test)]
1585mod tests {
1586 use super::*;
1587 use crate::column_chunk::NODE_GROUP_SIZE;
1588
1589 #[test]
1592 fn test_node_table_empty() {
1593 let table = NodeTable::new(
1594 1,
1595 "Person".into(),
1596 vec![
1597 ColumnDefinition {
1598 compression: akar_common::enums::CompressionType::Uncompressed,
1599 name: "name".into(),
1600 logical_type: LogicalTypeID::String,
1601 is_primary_key: true,
1602 },
1603 ColumnDefinition {
1604 compression: akar_common::enums::CompressionType::Uncompressed,
1605 name: "age".into(),
1606 logical_type: LogicalTypeID::Int64,
1607 is_primary_key: false,
1608 },
1609 ],
1610 );
1611 assert_eq!(table.num_rows, 0);
1612 assert!(table.node_groups.is_empty());
1613 }
1614
1615 #[test]
1616 fn test_node_table_insert_and_get() {
1617 let mut table = NodeTable::new(
1618 1,
1619 "Person".into(),
1620 vec![
1621 ColumnDefinition {
1622 compression: akar_common::enums::CompressionType::Uncompressed,
1623 name: "name".into(),
1624 logical_type: LogicalTypeID::String,
1625 is_primary_key: true,
1626 },
1627 ColumnDefinition {
1628 compression: akar_common::enums::CompressionType::Uncompressed,
1629 name: "age".into(),
1630 logical_type: LogicalTypeID::Int64,
1631 is_primary_key: false,
1632 },
1633 ],
1634 );
1635 table
1636 .insert_row(vec![Value::String("Alice".into()), Value::Int64(30)])
1637 .unwrap();
1638 table
1639 .insert_row(vec![Value::String("Bob".into()), Value::Int64(25)])
1640 .unwrap();
1641
1642 assert_eq!(table.num_rows, 2);
1643 assert_eq!(table.get_value(0, 0), Some(&Value::String("Alice".into())));
1644 assert_eq!(table.get_value(1, 1), Some(&Value::Int64(25)));
1645 }
1646
1647 #[test]
1648 fn test_delete_row_removes_pk_from_hash_and_art_index() {
1649 let mut table = NodeTable::new(
1650 1,
1651 "Person".into(),
1652 vec![
1653 ColumnDefinition {
1654 compression: akar_common::enums::CompressionType::Uncompressed,
1655 name: "name".into(),
1656 logical_type: LogicalTypeID::String,
1657 is_primary_key: true,
1658 },
1659 ColumnDefinition {
1660 compression: akar_common::enums::CompressionType::Uncompressed,
1661 name: "age".into(),
1662 logical_type: LogicalTypeID::Int64,
1663 is_primary_key: false,
1664 },
1665 ],
1666 );
1667 table.art_index = Some(ArtPrimaryKeyIndex::new("test_art"));
1668 table
1669 .insert_row(vec![Value::String("Alice".into()), Value::Int64(30)])
1670 .unwrap();
1671 table
1672 .insert_row(vec![Value::String("Bob".into()), Value::Int64(25)])
1673 .unwrap();
1674 let art = table.art_index.as_ref().unwrap();
1675 assert_eq!(
1676 art.lookup(&ArtKey::from_value(&Value::String("Alice".into())).unwrap()),
1677 Some(0)
1678 );
1679
1680 table.delete_row(0).unwrap();
1681
1682 assert!(table.lookup_by_pk(&Value::String("Alice".into())).is_none());
1684 let art = table.art_index.as_ref().unwrap();
1685 assert_eq!(art.len(), 1, "ART must drop the deleted entry");
1686 assert!(
1687 art.lookup(&ArtKey::from_value(&Value::String("Alice".into())).unwrap())
1688 .is_none()
1689 );
1690 assert!(
1691 art.lookup(&ArtKey::from_value(&Value::String("Bob".into())).unwrap())
1692 .is_some()
1693 );
1694 let hits = table.lookup_by_pk_range(
1696 Some(&Value::String("A".into())),
1697 true,
1698 Some(&Value::String("C".into())),
1699 true,
1700 100,
1701 );
1702 assert_eq!(hits, vec![1], "only 'Bob' (row 1) should be in range");
1703
1704 table
1706 .insert_row(vec![Value::String("Alice".into()), Value::Int64(31)])
1707 .unwrap();
1708 assert_eq!(table.lookup_by_pk(&Value::String("Alice".into())), Some(2));
1709 }
1710
1711 #[test]
1712 fn test_node_table_scan_column() {
1713 let mut table = NodeTable::new(
1714 1,
1715 "T".into(),
1716 vec![ColumnDefinition {
1717 compression: akar_common::enums::CompressionType::Uncompressed,
1718 name: "val".into(),
1719 logical_type: LogicalTypeID::Int64,
1720 is_primary_key: false,
1721 }],
1722 );
1723 for i in 0..100 {
1724 table.insert_row(vec![Value::Int64(i)]).unwrap();
1725 }
1726 let scanned = table.scan_column(0, 10, 5, None, &[]);
1727 assert_eq!(scanned.len(), 5);
1728 assert_eq!(scanned[0], Value::Int64(10));
1729 assert_eq!(scanned[4], Value::Int64(14));
1730 }
1731
1732 #[test]
1733 fn test_node_table_to_column_major() {
1734 let mut table = NodeTable::new(
1735 1,
1736 "T".into(),
1737 vec![
1738 ColumnDefinition {
1739 compression: akar_common::enums::CompressionType::Uncompressed,
1740 name: "x".into(),
1741 logical_type: LogicalTypeID::Int64,
1742 is_primary_key: false,
1743 },
1744 ColumnDefinition {
1745 compression: akar_common::enums::CompressionType::Uncompressed,
1746 name: "y".into(),
1747 logical_type: LogicalTypeID::Int64,
1748 is_primary_key: false,
1749 },
1750 ],
1751 );
1752 table.insert_row(vec![Value::Int64(1), Value::Int64(10)]).unwrap();
1753 table.insert_row(vec![Value::Int64(2), Value::Int64(20)]).unwrap();
1754
1755 let data = table.to_column_major_data();
1756 assert_eq!(data.len(), 2);
1757 assert_eq!(data[0], vec![Value::Int64(1), Value::Int64(2)]);
1758 assert_eq!(data[1], vec![Value::Int64(10), Value::Int64(20)]);
1759 }
1760
1761 #[test]
1762 fn test_node_table_auto_node_group() {
1763 let mut table = NodeTable::new(
1764 1,
1765 "T".into(),
1766 vec![ColumnDefinition {
1767 compression: akar_common::enums::CompressionType::Uncompressed,
1768 name: "v".into(),
1769 logical_type: LogicalTypeID::Int64,
1770 is_primary_key: false,
1771 }],
1772 );
1773 for i in 0..NODE_GROUP_SIZE as u64 + 1 {
1775 table.insert_row(vec![Value::Int64(i as i64)]).unwrap();
1776 }
1777 assert_eq!(table.num_rows, NODE_GROUP_SIZE as u64 + 1);
1778 assert_eq!(table.node_groups.len(), 2);
1779 assert_eq!(table.node_groups[0].num_nodes, NODE_GROUP_SIZE as u64);
1780 assert_eq!(table.node_groups[1].num_nodes, 1);
1781 assert_eq!(table.get_value(0, 0), Some(&Value::Int64(0)));
1783 assert_eq!(
1784 table.get_value(NODE_GROUP_SIZE, 0),
1785 Some(&Value::Int64(NODE_GROUP_SIZE as i64))
1786 );
1787 }
1788
1789 fn make_rel_table() -> RelTable {
1792 RelTable::new(
1793 1,
1794 "Knows".into(),
1795 0,
1796 1,
1797 vec![
1798 ColumnDefinition {
1799 compression: akar_common::enums::CompressionType::Uncompressed,
1800 name: "since".into(),
1801 logical_type: LogicalTypeID::Int64,
1802 is_primary_key: false,
1803 },
1804 ColumnDefinition {
1805 compression: akar_common::enums::CompressionType::Uncompressed,
1806 name: "weight".into(),
1807 logical_type: LogicalTypeID::Double,
1808 is_primary_key: false,
1809 },
1810 ],
1811 )
1812 }
1813
1814 #[test]
1815 fn test_rel_table_empty() {
1816 let rel = make_rel_table();
1817 assert_eq!(rel.num_rows, 0);
1818 assert!(rel.edges.is_empty());
1819 assert!(rel.fwd_adj.is_empty());
1820 assert!(rel.rev_adj.is_empty());
1821 }
1822
1823 #[test]
1824 fn test_rel_insert_basic() {
1825 let mut rel = make_rel_table();
1826 rel.insert_rel(0, 1, vec![Value::Int64(2020), Value::Double(0.5)])
1827 .unwrap();
1828 rel.insert_rel(0, 2, vec![Value::Int64(2021), Value::Double(0.8)])
1829 .unwrap();
1830 rel.insert_rel(1, 0, vec![Value::Int64(2020), Value::Double(0.3)])
1831 .unwrap();
1832
1833 assert_eq!(rel.num_rows, 3);
1834 assert_eq!(rel.edges.len(), 3);
1835
1836 let fwd = rel.scan_adj_list(0);
1838 assert_eq!(fwd.len(), 2);
1839 assert_eq!(fwd[0], (1, 0)); assert_eq!(fwd[1], (2, 1)); let fwd1 = rel.scan_adj_list(1);
1844 assert_eq!(fwd1.len(), 1);
1845 assert_eq!(fwd1[0], (0, 2));
1846 }
1847
1848 #[test]
1849 fn test_rel_reverse_adjacency() {
1850 let mut rel = make_rel_table();
1851 rel.insert_rel(0, 5, vec![Value::Int64(2022), Value::Double(1.0)])
1852 .unwrap();
1853 rel.insert_rel(3, 5, vec![Value::Int64(2023), Value::Double(1.5)])
1854 .unwrap();
1855
1856 let rev = rel.scan_rev_adj_list(5);
1858 assert_eq!(rev.len(), 2);
1859 assert_eq!(rev[0], (0, 0));
1860 assert_eq!(rev[1], (3, 1));
1861 }
1862
1863 #[test]
1864 fn test_rel_get_edge_properties() {
1865 let mut rel = make_rel_table();
1866 rel.insert_rel(0, 1, vec![Value::Int64(2020), Value::Double(0.5)])
1867 .unwrap();
1868 rel.insert_rel(2, 3, vec![Value::Int64(2021), Value::Double(0.9)])
1869 .unwrap();
1870
1871 let props0 = rel.get_edge_properties(0);
1872 assert_eq!(props0, vec![Value::Int64(2020), Value::Double(0.5)]);
1873
1874 let props1 = rel.get_edge_properties(1);
1875 assert_eq!(props1, vec![Value::Int64(2021), Value::Double(0.9)]);
1876 }
1877
1878 #[test]
1879 fn test_rel_get_outgoing_edges() {
1880 let mut rel = make_rel_table();
1881 rel.insert_rel(0, 10, vec![Value::Int64(2020), Value::Double(1.0)])
1882 .unwrap();
1883 rel.insert_rel(0, 20, vec![Value::Int64(2021), Value::Double(2.0)])
1884 .unwrap();
1885
1886 let outgoing = rel.get_outgoing_edges(0);
1887 assert_eq!(outgoing.len(), 2);
1888 assert_eq!(outgoing[0].0, 10);
1889 assert_eq!(outgoing[0].1, vec![Value::Int64(2020), Value::Double(1.0)]);
1890 assert_eq!(outgoing[1].0, 20);
1891 }
1892
1893 #[test]
1894 fn test_rel_get_incoming_edges() {
1895 let mut rel = make_rel_table();
1896 rel.insert_rel(10, 5, vec![Value::Int64(2020), Value::Double(1.0)])
1897 .unwrap();
1898 rel.insert_rel(20, 5, vec![Value::Int64(2021), Value::Double(2.0)])
1899 .unwrap();
1900
1901 let incoming = rel.get_incoming_edges(5);
1902 assert_eq!(incoming.len(), 2);
1903 assert_eq!(incoming[0].0, 10);
1904 assert_eq!(incoming[1].0, 20);
1905 }
1906
1907 #[test]
1908 fn test_rel_no_edges() {
1909 let rel = make_rel_table();
1910 assert!(rel.scan_adj_list(0).is_empty());
1911 assert!(rel.scan_rev_adj_list(0).is_empty());
1912 assert!(rel.get_outgoing_edges(0).is_empty());
1913 assert!(rel.get_incoming_edges(0).is_empty());
1914 }
1915
1916 #[test]
1917 fn test_rel_insert_row_legacy() {
1918 let mut rel = make_rel_table();
1919 rel.insert_row(vec![Value::Int64(2022), Value::Double(3.0)]).unwrap();
1921 assert_eq!(rel.num_rows, 1);
1922 assert_eq!(rel.edges[0], (0, 0)); assert_eq!(rel.get_edge_properties(0), vec![Value::Int64(2022), Value::Double(3.0)]);
1924 }
1925
1926 #[test]
1927 fn test_rel_wrong_column_count() {
1928 let mut rel = make_rel_table();
1929 let result = rel.insert_rel(0, 1, vec![Value::Int64(42)]); assert!(result.is_err());
1931 }
1932
1933 #[test]
1934 fn test_rel_get_column() {
1935 let mut rel = make_rel_table();
1936 rel.insert_rel(0, 1, vec![Value::Int64(2020), Value::Double(1.5)])
1937 .unwrap();
1938 rel.insert_rel(1, 2, vec![Value::Int64(2021), Value::Double(2.5)])
1939 .unwrap();
1940
1941 let since_col = rel.get_column(0).unwrap();
1942 assert_eq!(since_col, &[Value::Int64(2020), Value::Int64(2021)]);
1943
1944 let weight_col = rel.get_column(1).unwrap();
1945 assert_eq!(weight_col, &[Value::Double(1.5), Value::Double(2.5)]);
1946 }
1947
1948 #[test]
1949 fn test_rel_to_column_major() {
1950 let mut rel = make_rel_table();
1951 rel.insert_rel(0, 1, vec![Value::Int64(2020), Value::Double(0.5)])
1952 .unwrap();
1953 rel.insert_rel(2, 3, vec![Value::Int64(2021), Value::Double(0.9)])
1954 .unwrap();
1955
1956 let data = rel.to_column_major_data();
1957 assert_eq!(data.len(), 2);
1958 assert_eq!(data[0], vec![Value::Int64(2020), Value::Int64(2021)]);
1959 assert_eq!(data[1], vec![Value::Double(0.5), Value::Double(0.9)]);
1960 }
1961
1962 #[test]
1965 fn test_catalog_create_and_lookup() {
1966 let cat = TableCatalog::new();
1967 let node_table = cat.create_node_table(
1968 "Person".into(),
1969 vec![ColumnDefinition {
1970 compression: akar_common::enums::CompressionType::Uncompressed,
1971 name: "id".into(),
1972 logical_type: LogicalTypeID::Int64,
1973 is_primary_key: true,
1974 }],
1975 );
1976 assert_eq!(node_table.table_id, 0);
1977
1978 let rel_table = cat.create_rel_table(
1979 "Knows".into(),
1980 0,
1981 1,
1982 vec![ColumnDefinition {
1983 compression: akar_common::enums::CompressionType::Uncompressed,
1984 name: "since".into(),
1985 logical_type: LogicalTypeID::Int64,
1986 is_primary_key: false,
1987 }],
1988 );
1989 assert_eq!(rel_table.table_id, 1);
1990
1991 assert!(cat.get_node_table(0).is_some());
1992 assert!(cat.get_rel_table(1).is_some());
1993 assert_eq!(cat.node_table_num_rows("Person"), 0);
1994 }
1995
1996 #[test]
1997 fn test_refresh_vector_index_after_dml() {
1998 let cat = TableCatalog::new();
2002 cat.create_node_table(
2003 "Item".into(),
2004 vec![
2005 ColumnDefinition {
2006 compression: akar_common::enums::CompressionType::Uncompressed,
2007 name: "id".into(),
2008 logical_type: LogicalTypeID::Int64,
2009 is_primary_key: true,
2010 },
2011 ColumnDefinition {
2012 compression: akar_common::enums::CompressionType::Uncompressed,
2013 name: "embedding".into(),
2014 logical_type: LogicalTypeID::List,
2015 is_primary_key: false,
2016 },
2017 ],
2018 );
2019
2020 let vec3 = |x: f64, y: f64, z: f64| Value::List(vec![Value::Double(x), Value::Double(y), Value::Double(z)]);
2021 {
2022 let mut t = cat.get_node_table_by_name_mut("Item").unwrap();
2023 t.insert_row(vec![Value::Int64(1), vec3(1.0, 0.0, 0.0)]).unwrap();
2024 t.insert_row(vec![Value::Int64(2), vec3(0.0, 1.0, 0.0)]).unwrap();
2025 t.insert_row(vec![Value::Int64(3), vec3(0.0, 0.0, 1.0)]).unwrap();
2026 }
2027
2028 cat.create_vector_index(
2029 "item_vec".into(),
2030 "Item".into(),
2031 "embedding".into(),
2032 DistanceMetric::Cosine,
2033 3,
2034 );
2035
2036 cat.refresh_vector_indexes_for_tables(&[0]);
2038 {
2039 let vi = cat.get_vector_index(1).unwrap();
2040 assert_eq!(vi.hnsw().len(), 3);
2041 assert!(vi.hnsw().get_vector(2).is_some());
2042 }
2043
2044 {
2046 let mut t = cat.get_node_table_by_name_mut("Item").unwrap();
2047 t.insert_row(vec![Value::Int64(4), vec3(1.0, 1.0, 0.0)]).unwrap();
2048 t.insert_row(vec![Value::Int64(5), vec3(0.0, 1.0, 1.0)]).unwrap();
2049 }
2050
2051 cat.refresh_vector_indexes_for_tables(&[0]);
2052 let vi = cat.get_vector_index(1).unwrap();
2053 assert_eq!(vi.hnsw().len(), 5);
2054 assert!(vi.hnsw().get_vector(4).is_some());
2055 }
2056}