1use crate::runtime::wal::{Mutation, WriteAheadLog};
5use anyhow::Result;
6use std::collections::{HashMap, HashSet};
7use std::sync::Arc;
8use std::time::{SystemTime, UNIX_EPOCH};
9use tracing::{instrument, trace};
10use uni_common::core::id::{Eid, Vid};
11use uni_common::graph::simple_graph::{Direction, SimpleGraph};
12use uni_common::{Properties, Value};
13use uni_crdt::Crdt;
14
15#[derive(Debug, Default)]
23pub struct OccReadSet {
24 pub vertices: HashSet<Vid>,
26 pub edges: HashSet<Eid>,
28}
29
30impl OccReadSet {
31 pub fn is_empty(&self) -> bool {
34 self.vertices.is_empty() && self.edges.is_empty()
35 }
36}
37
38fn now_nanos() -> i64 {
40 SystemTime::now()
41 .duration_since(UNIX_EPOCH)
42 .map(|d| d.as_nanos() as i64)
43 .unwrap_or(0)
44}
45
46pub(crate) fn try_as_crdt(v: &Value) -> Option<Crdt> {
58 if !matches!(v, Value::Map(_)) {
59 return None;
60 }
61 serde_json::from_value::<Crdt>(v.clone().into()).ok()
62}
63
64pub fn serialize_constraint_key(label: &str, key_values: &[(String, Value)]) -> Vec<u8> {
67 let mut buf = label.as_bytes().to_vec();
68 buf.push(0); let mut sorted = key_values.to_vec();
70 sorted.sort_by(|a, b| a.0.cmp(&b.0));
71 for (k, v) in &sorted {
72 buf.extend(k.as_bytes());
73 buf.push(0);
74 buf.extend(serde_json::to_vec(v).unwrap_or_default());
76 buf.push(0);
77 }
78 buf
79}
80
81#[derive(Debug, Clone, Default)]
87pub struct MutationStats {
88 pub nodes_created: usize,
89 pub nodes_deleted: usize,
90 pub relationships_created: usize,
91 pub relationships_deleted: usize,
92 pub properties_set: usize,
93 pub properties_removed: usize,
94 pub labels_added: usize,
95 pub labels_removed: usize,
96}
97
98impl MutationStats {
99 pub fn diff(&self, before: &Self) -> Self {
101 Self {
102 nodes_created: self.nodes_created.saturating_sub(before.nodes_created),
103 nodes_deleted: self.nodes_deleted.saturating_sub(before.nodes_deleted),
104 relationships_created: self
105 .relationships_created
106 .saturating_sub(before.relationships_created),
107 relationships_deleted: self
108 .relationships_deleted
109 .saturating_sub(before.relationships_deleted),
110 properties_set: self.properties_set.saturating_sub(before.properties_set),
111 properties_removed: self
112 .properties_removed
113 .saturating_sub(before.properties_removed),
114 labels_added: self.labels_added.saturating_sub(before.labels_added),
115 labels_removed: self.labels_removed.saturating_sub(before.labels_removed),
116 }
117 }
118}
119
120#[derive(Clone, Debug)]
121pub struct TombstoneEntry {
122 pub eid: Eid,
123 pub src_vid: Vid,
124 pub dst_vid: Vid,
125 pub edge_type: u32,
126}
127
128pub struct L0Buffer {
129 pub graph: SimpleGraph,
131 pub tombstones: HashMap<Eid, TombstoneEntry>,
133 pub vertex_tombstones: HashSet<Vid>,
135 pub edge_versions: HashMap<Eid, u64>,
137 pub vertex_versions: HashMap<Vid, u64>,
139 pub edge_properties: HashMap<Eid, Properties>,
141 pub vertex_properties: HashMap<Vid, Properties>,
143 pub edge_endpoints: HashMap<Eid, (Vid, Vid, u32)>,
145 pub vertex_labels: HashMap<Vid, Vec<String>>,
148 pub label_to_vids: HashMap<String, HashSet<Vid>>,
151 pub vertex_label_overwrites: HashSet<Vid>,
159 pub edge_types: HashMap<Eid, String>,
161 pub current_version: u64,
163 pub mutation_count: usize,
165 pub mutation_stats: MutationStats,
167 pub wal: Option<Arc<WriteAheadLog>>,
169 pub wal_lsn_at_flush: u64,
172 pub wal_lsn_at_start: u64,
182 pub vertex_created_at: HashMap<Vid, i64>,
184 pub vertex_updated_at: HashMap<Vid, i64>,
186 pub edge_created_at: HashMap<Eid, i64>,
188 pub edge_updated_at: HashMap<Eid, i64>,
190 pub estimated_size: usize,
193 pub constraint_index: HashMap<Vec<u8>, Vid>,
197 pub merge_guard_index: HashMap<Vec<u8>, Vid>,
207 pub edge_constraint_index: HashMap<Vec<u8>, Eid>,
214 pub extid_index: HashMap<String, Vid>,
221 pub vertex_partial_keys: HashMap<Vid, HashSet<String>>,
227 pub edge_partial_keys: HashMap<Eid, HashSet<String>>,
234 pub pending_embeddings: HashMap<Vid, String>,
241 pub occ_read_seq: u64,
246 pub occ_read_set: Option<Arc<parking_lot::Mutex<OccReadSet>>>,
251 pub plugin_registry: Option<Arc<uni_plugin::PluginRegistry>>,
260}
261
262impl std::fmt::Debug for L0Buffer {
263 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
264 f.debug_struct("L0Buffer")
265 .field("vertex_count", &self.graph.vertex_count())
266 .field("edge_count", &self.graph.edge_count())
267 .field("tombstones", &self.tombstones.len())
268 .field("vertex_tombstones", &self.vertex_tombstones.len())
269 .field("current_version", &self.current_version)
270 .field("mutation_count", &self.mutation_count)
271 .finish()
272 }
273}
274
275impl Clone for L0Buffer {
276 fn clone(&self) -> Self {
281 Self {
282 graph: self.graph.clone(),
283 tombstones: self.tombstones.clone(),
284 vertex_tombstones: self.vertex_tombstones.clone(),
285 edge_versions: self.edge_versions.clone(),
286 vertex_versions: self.vertex_versions.clone(),
287 edge_properties: self.edge_properties.clone(),
288 vertex_properties: self.vertex_properties.clone(),
289 edge_endpoints: self.edge_endpoints.clone(),
290 vertex_labels: self.vertex_labels.clone(),
291 label_to_vids: self.label_to_vids.clone(),
292 vertex_label_overwrites: self.vertex_label_overwrites.clone(),
293 edge_types: self.edge_types.clone(),
294 current_version: self.current_version,
295 mutation_count: self.mutation_count,
296 mutation_stats: self.mutation_stats.clone(),
297 wal: None, wal_lsn_at_flush: self.wal_lsn_at_flush,
299 wal_lsn_at_start: self.wal_lsn_at_start,
300 vertex_created_at: self.vertex_created_at.clone(),
301 vertex_updated_at: self.vertex_updated_at.clone(),
302 edge_created_at: self.edge_created_at.clone(),
303 edge_updated_at: self.edge_updated_at.clone(),
304 estimated_size: self.estimated_size,
305 constraint_index: self.constraint_index.clone(),
306 edge_constraint_index: self.edge_constraint_index.clone(),
307 merge_guard_index: self.merge_guard_index.clone(),
308 extid_index: self.extid_index.clone(),
309 vertex_partial_keys: self.vertex_partial_keys.clone(),
310 edge_partial_keys: self.edge_partial_keys.clone(),
311 pending_embeddings: self.pending_embeddings.clone(),
312 occ_read_seq: self.occ_read_seq,
313 occ_read_set: None,
315 plugin_registry: self.plugin_registry.clone(),
318 }
319 }
320}
321
322impl L0Buffer {
323 fn append_unique_labels(existing: &mut Vec<String>, labels: &[String]) {
325 for label in labels {
326 if !existing.contains(label) {
327 existing.push(label.clone());
328 }
329 }
330 }
331
332 fn index_labels_for_vid(&mut self, vid: Vid, labels: &[String]) {
334 for label in labels {
335 self.label_to_vids
336 .entry(label.clone())
337 .or_default()
338 .insert(vid);
339 }
340 }
341
342 fn extid_of(props: &Properties) -> Option<String> {
344 props
345 .get("ext_id")
346 .and_then(|v| v.as_str())
347 .map(str::to_owned)
348 }
349
350 fn sync_extid_index(&mut self, vid: Vid, old: Option<String>, new: Option<String>) {
358 if old == new {
359 return;
360 }
361 if let Some(old) = old
362 && self.extid_index.get(&old) == Some(&vid)
363 {
364 self.extid_index.remove(&old);
365 }
366 if let Some(new) = new {
367 self.extid_index.insert(new, vid);
368 }
369 }
370
371 fn remove_vid_from_label_index(&mut self, vid: Vid) {
373 if let Some(labels) = self.vertex_labels.get(&vid) {
374 for label in labels {
375 if let Some(set) = self.label_to_vids.get_mut(label) {
376 set.remove(&vid);
377 }
378 }
379 }
380 }
381
382 pub fn set_vertex_labels(&mut self, vid: Vid, labels: &[String]) {
393 self.remove_vid_from_label_index(vid);
394 self.vertex_labels.insert(vid, labels.to_vec());
395 self.index_labels_for_vid(vid, labels);
396 self.vertex_label_overwrites.insert(vid);
397 self.current_version += 1;
398 self.mutation_count += 1;
399 }
400
401 fn merge_crdt_properties(
410 entry: &mut Properties,
411 properties: Properties,
412 registry: Option<&Arc<uni_plugin::PluginRegistry>>,
413 ) {
414 if entry.is_empty() {
416 *entry = properties;
417 return;
418 }
419
420 for (k, v) in properties {
421 if let Some(mut new_crdt) = try_as_crdt(&v)
425 && let Some(existing_v) = entry.get(&k)
426 && let Ok(existing_crdt) = serde_json::from_value::<Crdt>(existing_v.clone().into())
427 {
428 let merged = match registry {
433 Some(reg) => new_crdt.merge_via_registry(&existing_crdt, reg).is_ok(),
434 None => new_crdt.try_merge(&existing_crdt).is_ok(),
435 };
436 if merged && let Ok(merged_json) = serde_json::to_value(new_crdt) {
437 entry.insert(k, uni_common::Value::from(merged_json));
438 continue;
439 }
440 tracing::warn!(
448 property = %k,
449 existing_variant = existing_crdt.type_name(),
450 "overwriting CRDT property with a different CRDT variant \
451 (last-writer-wins); merged CRDT state is discarded"
452 );
453 } else if try_as_crdt(&v).is_none()
454 && entry.get(&k).is_some_and(|e| try_as_crdt(e).is_some())
455 {
456 tracing::warn!(
463 property = %k,
464 "overwriting CRDT property with non-CRDT value (last-writer-wins); \
465 merged CRDT state is discarded"
466 );
467 }
468 entry.insert(k, v);
470 }
471 }
472
473 fn estimate_properties_size(props: &Properties) -> usize {
475 props.keys().map(|k| k.len() + 32).sum()
476 }
477
478 pub fn size_bytes(&self) -> usize {
481 let mut total = 0;
482
483 total += self.graph.vertex_count() * 8;
485 total += self.graph.edge_count() * 24;
486
487 for props in self.vertex_properties.values() {
489 total += Self::estimate_properties_size(props);
490 }
491 for props in self.edge_properties.values() {
492 total += Self::estimate_properties_size(props);
493 }
494
495 total += self.tombstones.len() * 64;
497 total += self.vertex_tombstones.len() * 8;
498 total += self.edge_versions.len() * 16;
499 total += self.vertex_versions.len() * 16;
500 total += self.edge_endpoints.len() * 28; for labels in self.vertex_labels.values() {
504 total += labels.iter().map(|l| l.len() + 24).sum::<usize>();
505 }
506
507 for (label, vids) in &self.label_to_vids {
509 total += label.len() + 24 + vids.len() * 8 + 48; }
511
512 for type_name in self.edge_types.values() {
514 total += type_name.len() + 24;
515 }
516
517 total += self.vertex_created_at.len() * 16;
519 total += self.vertex_updated_at.len() * 16;
520 total += self.edge_created_at.len() * 16;
521 total += self.edge_updated_at.len() * 16;
522
523 total
524 }
525
526 pub fn new(start_version: u64, wal: Option<Arc<WriteAheadLog>>) -> Self {
527 Self {
528 graph: SimpleGraph::new(),
529 tombstones: HashMap::new(),
530 vertex_tombstones: HashSet::new(),
531 edge_versions: HashMap::new(),
532 vertex_versions: HashMap::new(),
533 edge_properties: HashMap::new(),
534 vertex_properties: HashMap::new(),
535 edge_endpoints: HashMap::new(),
536 vertex_labels: HashMap::new(),
537 label_to_vids: HashMap::new(),
538 vertex_label_overwrites: HashSet::new(),
539 edge_types: HashMap::new(),
540 current_version: start_version,
541 mutation_count: 0,
542 mutation_stats: MutationStats::default(),
543 wal,
544 wal_lsn_at_flush: 0,
545 wal_lsn_at_start: 0,
546 vertex_created_at: HashMap::new(),
547 vertex_updated_at: HashMap::new(),
548 edge_created_at: HashMap::new(),
549 edge_updated_at: HashMap::new(),
550 estimated_size: 0,
551 constraint_index: HashMap::new(),
552 edge_constraint_index: HashMap::new(),
553 merge_guard_index: HashMap::new(),
554 extid_index: HashMap::new(),
555 vertex_partial_keys: HashMap::new(),
556 edge_partial_keys: HashMap::new(),
557 pending_embeddings: HashMap::new(),
558 occ_read_seq: 0,
559 occ_read_set: None,
560 plugin_registry: None,
561 }
562 }
563
564 pub fn set_plugin_registry(&mut self, registry: Arc<uni_plugin::PluginRegistry>) {
571 self.plugin_registry = Some(registry);
572 }
573
574 pub fn insert_vertex(&mut self, vid: Vid, properties: Properties) {
575 self.insert_vertex_with_labels(vid, properties, &[]);
576 }
577
578 pub fn insert_vertex_with_labels(
580 &mut self,
581 vid: Vid,
582 properties: Properties,
583 labels: &[String],
584 ) {
585 self.insert_vertex_with_labels_impl(vid, properties, labels, false);
586 }
587
588 fn insert_vertex_with_labels_impl(
591 &mut self,
592 vid: Vid,
593 properties: Properties,
594 labels: &[String],
595 skip_wal: bool,
596 ) {
597 self.current_version += 1;
598 let version = self.current_version;
599 let now = now_nanos();
600
601 if !skip_wal && let Some(wal) = &self.wal {
602 let _ = wal.append(Mutation::InsertVertex {
603 vid,
604 properties: properties.clone(),
605 labels: labels.to_vec(),
606 });
607 }
608
609 self.vertex_partial_keys.remove(&vid);
612
613 self.apply_vertex_write(vid, properties, labels, version, now);
614 self.mutation_stats.nodes_created += 1;
615 }
616
617 pub fn insert_vertex_partial_full(
638 &mut self,
639 vid: Vid,
640 props: Properties,
641 touched_keys: HashSet<String>,
642 labels: &[String],
643 ) {
644 self.insert_vertex_with_labels_partial_impl(vid, props, labels, false);
648 self.vertex_partial_keys
649 .entry(vid)
650 .or_default()
651 .extend(touched_keys);
652 }
653
654 pub fn insert_vertex_partial(&mut self, vid: Vid, touched: Properties, labels: &[String]) {
658 let touched_keys: Vec<String> = touched.keys().cloned().collect();
662
663 let already_full = self.vertex_properties.contains_key(&vid)
674 && !self.vertex_partial_keys.contains_key(&vid);
675
676 self.insert_vertex_with_labels_partial_impl(vid, touched, labels, false);
681
682 if !already_full {
683 self.vertex_partial_keys
684 .entry(vid)
685 .or_default()
686 .extend(touched_keys);
687 }
688 }
689
690 fn insert_vertex_with_labels_partial_impl(
694 &mut self,
695 vid: Vid,
696 properties: Properties,
697 labels: &[String],
698 skip_wal: bool,
699 ) {
700 self.current_version += 1;
701 let version = self.current_version;
702 let now = now_nanos();
703
704 if !skip_wal && let Some(wal) = &self.wal {
705 let _ = wal.append(Mutation::InsertVertex {
711 vid,
712 properties: properties.clone(),
713 labels: labels.to_vec(),
714 });
715 }
716
717 self.apply_vertex_write(vid, properties, labels, version, now);
722 }
723
724 fn apply_vertex_write(
732 &mut self,
733 vid: Vid,
734 properties: Properties,
735 labels: &[String],
736 version: u64,
737 now: i64,
738 ) {
739 self.vertex_tombstones.remove(&vid);
740
741 let props_size = Self::estimate_properties_size(&properties);
744 let props_count = properties.len();
745 let tracks_extid = properties.contains_key("ext_id");
746
747 let entry = self.vertex_properties.entry(vid).or_default();
748 let old_extid = if tracks_extid {
749 Self::extid_of(entry)
750 } else {
751 None
752 };
753 Self::merge_crdt_properties(entry, properties, self.plugin_registry.as_ref());
754 if tracks_extid {
755 let new_extid =
756 Self::extid_of(self.vertex_properties.get(&vid).expect("just inserted"));
757 self.sync_extid_index(vid, old_extid, new_extid);
758 }
759 self.vertex_versions.insert(vid, version);
760
761 self.vertex_created_at.entry(vid).or_insert(now);
763 self.vertex_updated_at.insert(vid, now);
764
765 let labels_size: usize = labels.iter().map(|l| l.len() + 24).sum();
768 let existing = self.vertex_labels.entry(vid).or_default();
769 Self::append_unique_labels(existing, labels);
770 self.index_labels_for_vid(vid, labels);
771
772 self.graph.add_vertex(vid);
773 self.mutation_count += 1;
774 self.mutation_stats.properties_set += props_count;
775 self.mutation_stats.labels_added += labels.len();
776
777 self.estimated_size += 8 + props_size + 16 + labels_size + 32;
778 }
779
780 pub fn add_vertex_labels(&mut self, vid: Vid, labels: &[String]) {
782 let existing = self.vertex_labels.entry(vid).or_default();
783 Self::append_unique_labels(existing, labels);
784 self.index_labels_for_vid(vid, labels);
785 }
786
787 pub fn remove_vertex_label(&mut self, vid: Vid, label: &str) -> bool {
790 if let Some(labels) = self.vertex_labels.get_mut(&vid)
791 && let Some(pos) = labels.iter().position(|l| l == label)
792 {
793 labels.remove(pos);
794 if let Some(set) = self.label_to_vids.get_mut(label) {
795 set.remove(&vid);
796 }
797 self.current_version += 1;
798 self.mutation_count += 1;
799 self.mutation_stats.labels_removed += 1;
800 return true;
803 }
804 false
805 }
806
807 pub fn set_edge_type(&mut self, eid: Eid, edge_type: String) {
809 self.edge_types.insert(eid, edge_type);
810 }
811
812 pub fn delete_vertex(&mut self, vid: Vid) -> Result<()> {
813 self.delete_vertex_impl(vid, false)
814 }
815
816 fn delete_vertex_impl(&mut self, vid: Vid, skip_wal: bool) -> Result<()> {
819 self.current_version += 1;
820
821 if !skip_wal && let Some(wal) = &mut self.wal {
822 let labels = self.vertex_labels.get(&vid).cloned().unwrap_or_default();
823 wal.append(Mutation::DeleteVertex { vid, labels })?;
824 }
825
826 self.apply_vertex_deletion(vid);
827 Ok(())
828 }
829
830 fn apply_vertex_deletion(&mut self, vid: Vid) {
834 let version = self.current_version;
835
836 let mut edges_to_remove = HashSet::new();
838
839 for entry in self.graph.neighbors(vid, Direction::Outgoing) {
841 edges_to_remove.insert(entry.eid);
842 }
843
844 for entry in self.graph.neighbors(vid, Direction::Incoming) {
846 edges_to_remove.insert(entry.eid); }
848
849 let cascaded_edges_count = edges_to_remove.len();
850
851 for eid in edges_to_remove {
853 if let Some((src, dst, etype)) = self.edge_endpoints.get(&eid) {
855 self.tombstones.insert(
856 eid,
857 TombstoneEntry {
858 eid,
859 src_vid: *src,
860 dst_vid: *dst,
861 edge_type: *etype,
862 },
863 );
864 self.edge_versions.insert(eid, version);
865 self.edge_endpoints.remove(&eid);
866 self.edge_properties.remove(&eid);
867 self.graph.remove_edge(eid);
868 self.mutation_count += 1;
869 self.mutation_stats.relationships_deleted += 1;
870 }
871 }
872
873 self.remove_vid_from_label_index(vid);
874 self.vertex_tombstones.insert(vid);
875 if let Some(props) = self.vertex_properties.get(&vid)
878 && let Some(ext) = Self::extid_of(props)
879 && self.extid_index.get(&ext) == Some(&vid)
880 {
881 self.extid_index.remove(&ext);
882 }
883 self.vertex_properties.remove(&vid);
884 self.vertex_partial_keys.remove(&vid);
886 self.vertex_versions.insert(vid, version);
887 self.graph.remove_vertex(vid);
888 self.mutation_count += 1;
889 self.mutation_stats.nodes_deleted += 1;
890
891 self.constraint_index.retain(|_, v| *v != vid);
893 self.merge_guard_index.retain(|_, v| *v != vid);
896
897 self.estimated_size += cascaded_edges_count * 72 + 8;
899 }
900
901 pub fn insert_edge(
902 &mut self,
903 src_vid: Vid,
904 dst_vid: Vid,
905 edge_type: u32,
906 eid: Eid,
907 properties: Properties,
908 edge_type_name: Option<String>,
909 ) -> Result<()> {
910 self.insert_edge_impl(
911 src_vid,
912 dst_vid,
913 edge_type,
914 eid,
915 properties,
916 edge_type_name,
917 false,
918 )
919 }
920
921 #[allow(clippy::too_many_arguments)]
924 fn insert_edge_impl(
925 &mut self,
926 src_vid: Vid,
927 dst_vid: Vid,
928 edge_type: u32,
929 eid: Eid,
930 properties: Properties,
931 edge_type_name: Option<String>,
932 skip_wal: bool,
933 ) -> Result<()> {
934 self.current_version += 1;
935 let now = now_nanos();
936
937 if !skip_wal && let Some(wal) = &mut self.wal {
938 wal.append(Mutation::InsertEdge {
939 src_vid,
940 dst_vid,
941 edge_type,
942 eid,
943 version: self.current_version,
944 properties: properties.clone(),
945 edge_type_name: edge_type_name.clone(),
946 })?;
947 }
948
949 self.apply_edge_insertion(src_vid, dst_vid, edge_type, eid, properties)?;
950
951 let type_name_size = if let Some(ref name) = edge_type_name {
953 let size = name.len() + 24;
954 self.edge_types.insert(eid, name.clone());
955 size
956 } else {
957 0
958 };
959
960 self.edge_created_at.entry(eid).or_insert(now);
962 self.edge_updated_at.insert(eid, now);
963
964 self.edge_partial_keys.remove(&eid);
967
968 self.estimated_size += type_name_size;
969
970 Ok(())
971 }
972
973 #[allow(clippy::too_many_arguments)]
978 pub fn insert_edge_partial_full(
979 &mut self,
980 src_vid: Vid,
981 dst_vid: Vid,
982 edge_type: u32,
983 eid: Eid,
984 properties: Properties,
985 edge_type_name: Option<String>,
986 touched_keys: HashSet<String>,
987 ) -> Result<()> {
988 self.current_version += 1;
989 let now = now_nanos();
990
991 if let Some(wal) = &mut self.wal {
992 wal.append(Mutation::InsertEdge {
993 src_vid,
994 dst_vid,
995 edge_type,
996 eid,
997 version: self.current_version,
998 properties: properties.clone(),
999 edge_type_name: edge_type_name.clone(),
1000 })?;
1001 }
1002
1003 self.apply_edge_insertion(src_vid, dst_vid, edge_type, eid, properties)?;
1004
1005 self.edge_partial_keys
1009 .entry(eid)
1010 .or_default()
1011 .extend(touched_keys);
1012
1013 let type_name_size = if let Some(ref name) = edge_type_name {
1014 let size = name.len() + 24;
1015 self.edge_types.insert(eid, name.clone());
1016 size
1017 } else {
1018 0
1019 };
1020
1021 self.edge_created_at.entry(eid).or_insert(now);
1022 self.edge_updated_at.insert(eid, now);
1023
1024 self.estimated_size += type_name_size;
1025
1026 Ok(())
1027 }
1028
1029 fn apply_edge_insertion(
1038 &mut self,
1039 src_vid: Vid,
1040 dst_vid: Vid,
1041 edge_type: u32,
1042 eid: Eid,
1043 properties: Properties,
1044 ) -> Result<()> {
1045 let version = self.current_version;
1046
1047 if self.vertex_tombstones.contains(&src_vid) {
1050 anyhow::bail!(
1051 "Cannot insert edge: source vertex {} has been deleted (issue #77)",
1052 src_vid
1053 );
1054 }
1055 if self.vertex_tombstones.contains(&dst_vid) {
1056 anyhow::bail!(
1057 "Cannot insert edge: destination vertex {} has been deleted (issue #77)",
1058 dst_vid
1059 );
1060 }
1061
1062 if !self.graph.contains_vertex(src_vid) {
1067 self.graph.add_vertex(src_vid);
1068 }
1069 if !self.graph.contains_vertex(dst_vid) {
1070 self.graph.add_vertex(dst_vid);
1071 }
1072
1073 self.graph.add_edge(src_vid, dst_vid, eid, edge_type);
1074
1075 let props_size = Self::estimate_properties_size(&properties);
1077 let props_count = properties.len();
1078 if !properties.is_empty() {
1079 let entry = self.edge_properties.entry(eid).or_default();
1080 Self::merge_crdt_properties(entry, properties, self.plugin_registry.as_ref());
1081 }
1082
1083 self.edge_versions.insert(eid, version);
1084 self.edge_endpoints
1085 .insert(eid, (src_vid, dst_vid, edge_type));
1086 self.tombstones.remove(&eid);
1087 self.mutation_count += 1;
1088 self.mutation_stats.relationships_created += 1;
1089 self.mutation_stats.properties_set += props_count;
1090
1091 self.estimated_size += 24 + props_size + 16 + 28 + 32;
1093
1094 Ok(())
1095 }
1096
1097 pub fn delete_edge(
1098 &mut self,
1099 eid: Eid,
1100 src_vid: Vid,
1101 dst_vid: Vid,
1102 edge_type: u32,
1103 ) -> Result<()> {
1104 self.delete_edge_impl(eid, src_vid, dst_vid, edge_type, false)
1105 }
1106
1107 fn delete_edge_impl(
1110 &mut self,
1111 eid: Eid,
1112 src_vid: Vid,
1113 dst_vid: Vid,
1114 edge_type: u32,
1115 skip_wal: bool,
1116 ) -> Result<()> {
1117 self.current_version += 1;
1118 let now = now_nanos();
1119
1120 if !skip_wal && let Some(wal) = &mut self.wal {
1121 wal.append(Mutation::DeleteEdge {
1122 eid,
1123 src_vid,
1124 dst_vid,
1125 edge_type,
1126 version: self.current_version,
1127 })?;
1128 }
1129
1130 self.apply_edge_deletion(eid, src_vid, dst_vid, edge_type);
1131
1132 self.edge_updated_at.insert(eid, now);
1134
1135 Ok(())
1136 }
1137
1138 fn apply_edge_deletion(&mut self, eid: Eid, src_vid: Vid, dst_vid: Vid, edge_type: u32) {
1142 let version = self.current_version;
1143
1144 self.tombstones.insert(
1145 eid,
1146 TombstoneEntry {
1147 eid,
1148 src_vid,
1149 dst_vid,
1150 edge_type,
1151 },
1152 );
1153 self.edge_versions.insert(eid, version);
1154 self.edge_partial_keys.remove(&eid);
1157 self.edge_constraint_index.retain(|_, e| *e != eid);
1160 self.graph.remove_edge(eid);
1161 self.mutation_count += 1;
1162 self.mutation_stats.relationships_deleted += 1;
1163
1164 self.estimated_size += 80;
1166 }
1167
1168 pub fn get_neighbors(
1171 &self,
1172 vid: Vid,
1173 edge_type: u32,
1174 direction: Direction,
1175 ) -> Vec<(Vid, Eid, u64)> {
1176 let edges = self.graph.neighbors(vid, direction);
1177
1178 edges
1179 .iter()
1180 .filter(|e| e.edge_type == edge_type && !self.is_tombstoned(e.eid))
1181 .map(|e| {
1182 let neighbor = match direction {
1183 Direction::Outgoing => e.dst_vid,
1184 Direction::Incoming => e.src_vid,
1185 };
1186 let version = self.edge_versions.get(&e.eid).copied().unwrap_or(0);
1187 (neighbor, e.eid, version)
1188 })
1189 .collect()
1190 }
1191
1192 pub fn is_tombstoned(&self, eid: Eid) -> bool {
1193 self.tombstones.contains_key(&eid)
1194 }
1195
1196 pub fn vids_for_label(&self, label_name: &str) -> Vec<Vid> {
1199 self.label_to_vids
1200 .get(label_name)
1201 .map(|set| set.iter().copied().collect())
1202 .unwrap_or_default()
1203 }
1204
1205 pub fn all_vertex_vids(&self) -> Vec<Vid> {
1209 self.vertex_properties.keys().copied().collect()
1210 }
1211
1212 pub fn vids_for_labels(&self, label_names: &[&str]) -> Vec<Vid> {
1215 let mut result = HashSet::new();
1216 for label_name in label_names {
1217 if let Some(set) = self.label_to_vids.get(*label_name) {
1218 result.extend(set.iter().copied());
1219 }
1220 }
1221 result.into_iter().collect()
1222 }
1223
1224 pub fn vids_with_all_labels(&self, label_names: &[&str]) -> Vec<Vid> {
1227 if label_names.is_empty() {
1228 return Vec::new();
1229 }
1230 let sets: Vec<&HashSet<Vid>> = match label_names
1233 .iter()
1234 .map(|ln| self.label_to_vids.get(*ln))
1235 .collect::<Option<Vec<_>>>()
1236 {
1237 Some(s) => s,
1238 None => return Vec::new(),
1239 };
1240 let smallest = sets.iter().min_by_key(|s| s.len()).unwrap();
1242 smallest
1243 .iter()
1244 .copied()
1245 .filter(|vid| sets.iter().all(|s| s.contains(vid)))
1246 .collect()
1247 }
1248
1249 pub fn get_vertex_labels(&self, vid: Vid) -> Option<&[String]> {
1251 self.vertex_labels.get(&vid).map(|v| v.as_slice())
1252 }
1253
1254 pub fn get_edge_type(&self, eid: Eid) -> Option<&str> {
1256 self.edge_types.get(&eid).map(|s| s.as_str())
1257 }
1258
1259 pub fn eids_for_type(&self, type_name: &str) -> Vec<Eid> {
1262 self.edge_types
1263 .iter()
1264 .filter(|(eid, etype)| *etype == type_name && !self.tombstones.contains_key(eid))
1265 .map(|(eid, _)| *eid)
1266 .collect()
1267 }
1268
1269 pub fn all_edge_eids(&self) -> Vec<Eid> {
1273 self.edge_endpoints
1274 .keys()
1275 .filter(|eid| !self.tombstones.contains_key(eid))
1276 .copied()
1277 .collect()
1278 }
1279
1280 pub fn get_edge_endpoints(&self, eid: Eid) -> Option<(Vid, Vid)> {
1282 self.edge_endpoints
1283 .get(&eid)
1284 .map(|(src, dst, _)| (*src, *dst))
1285 }
1286
1287 pub fn get_edge_endpoint_full(&self, eid: Eid) -> Option<(Vid, Vid, u32)> {
1289 self.edge_endpoints.get(&eid).copied()
1290 }
1291
1292 pub fn insert_constraint_key(&mut self, key: Vec<u8>, vid: Vid) {
1294 self.constraint_index.insert(key, vid);
1295 }
1296
1297 pub fn has_constraint_key(&self, key: &[u8], exclude_vid: Vid) -> bool {
1300 self.constraint_index
1301 .get(key)
1302 .is_some_and(|&v| v != exclude_vid)
1303 }
1304
1305 pub fn insert_edge_constraint_key(&mut self, key: Vec<u8>, eid: Eid) {
1308 self.edge_constraint_index.insert(key, eid);
1309 }
1310
1311 pub fn has_edge_constraint_key(&self, key: &[u8], exclude_eid: Eid) -> bool {
1315 self.edge_constraint_index
1316 .get(key)
1317 .is_some_and(|&e| e != exclude_eid)
1318 }
1319
1320 pub fn insert_merge_guard_key(&mut self, key: Vec<u8>, vid: Vid) {
1322 self.merge_guard_index.insert(key, vid);
1323 }
1324
1325 pub fn has_merge_guard_key(&self, key: &[u8], exclude_vid: Vid) -> bool {
1328 self.merge_guard_index
1329 .get(key)
1330 .is_some_and(|&v| v != exclude_vid)
1331 }
1332
1333 #[instrument(skip(self, other), level = "trace")]
1334 pub fn validate_merge_edge_endpoints(&self, other: &L0Buffer) -> Result<()> {
1353 let is_deleted = |vid: &Vid| {
1357 (self.vertex_tombstones.contains(vid) || other.vertex_tombstones.contains(vid))
1358 && !other.vertex_properties.contains_key(vid)
1359 };
1360 for (eid, (src_vid, dst_vid, _etype)) in &other.edge_endpoints {
1361 if other.tombstones.contains_key(eid) {
1362 continue; }
1364 if is_deleted(src_vid) {
1365 anyhow::bail!(
1366 "Cannot insert edge {}: source vertex {} has been deleted (issue #77)",
1367 eid,
1368 src_vid
1369 );
1370 }
1371 if is_deleted(dst_vid) {
1372 anyhow::bail!(
1373 "Cannot insert edge {}: destination vertex {} has been deleted (issue #77)",
1374 eid,
1375 dst_vid
1376 );
1377 }
1378 }
1379 Ok(())
1380 }
1381
1382 pub fn merge(&mut self, other: &L0Buffer) -> Result<()> {
1383 self.validate_merge_edge_endpoints(other)?;
1387 self.merge_validated(
1388 other,
1389 other.vertex_properties.clone(),
1390 other.edge_properties.clone(),
1391 )
1392 }
1393
1394 pub fn merge_take(&mut self, other: &mut L0Buffer) -> Result<()> {
1403 self.validate_merge_edge_endpoints(other)?;
1406 let vertex_props = std::mem::take(&mut other.vertex_properties);
1407 let edge_props = std::mem::take(&mut other.edge_properties);
1408 self.merge_validated(other, vertex_props, edge_props)
1409 }
1410
1411 fn merge_validated(
1415 &mut self,
1416 other: &L0Buffer,
1417 vertex_props: HashMap<Vid, Properties>,
1418 mut edge_props: HashMap<Eid, Properties>,
1419 ) -> Result<()> {
1420 trace!(
1421 other_mutation_count = other.mutation_count,
1422 "Merging L0 buffer"
1423 );
1424 for &vid in &other.vertex_tombstones {
1429 self.delete_vertex_impl(vid, true)?;
1430 }
1431
1432 for (vid, props) in vertex_props {
1433 let labels = other.vertex_labels.get(&vid).cloned().unwrap_or_default();
1434 self.insert_vertex_with_labels_impl(vid, props, &labels, true);
1435 }
1436
1437 for (vid, labels) in &other.vertex_labels {
1439 if !self.vertex_labels.contains_key(vid) {
1440 self.vertex_labels.insert(*vid, labels.clone());
1441 for label in labels {
1442 self.label_to_vids
1443 .entry(label.clone())
1444 .or_default()
1445 .insert(*vid);
1446 }
1447 }
1448 }
1449
1450 for vid in &other.vertex_label_overwrites {
1457 if other.vertex_tombstones.contains(vid) {
1458 continue;
1459 }
1460 let labels = other.vertex_labels.get(vid).cloned().unwrap_or_default();
1461 self.remove_vid_from_label_index(*vid);
1462 self.vertex_labels.insert(*vid, labels.clone());
1463 self.index_labels_for_vid(*vid, &labels);
1464 self.vertex_label_overwrites.insert(*vid);
1469 }
1470
1471 for (eid, (src, dst, etype)) in &other.edge_endpoints {
1473 if other.tombstones.contains_key(eid) {
1474 self.delete_edge_impl(*eid, *src, *dst, *etype, true)?;
1475 } else {
1476 let props = edge_props.remove(eid).unwrap_or_default();
1477 let etype_name = other.edge_types.get(eid).cloned();
1478 self.insert_edge_impl(*src, *dst, *etype, *eid, props, etype_name, true)?;
1479 }
1480 }
1481
1482 for (eid, tombstone) in &other.tombstones {
1486 if !other.edge_endpoints.contains_key(eid) {
1487 self.delete_edge_impl(
1488 *eid,
1489 tombstone.src_vid,
1490 tombstone.dst_vid,
1491 tombstone.edge_type,
1492 true,
1493 )?;
1494 }
1495 }
1496
1497 for (vid, ts) in &other.vertex_created_at {
1502 self.vertex_created_at.entry(*vid).or_insert(*ts); }
1504 for (vid, ts) in &other.vertex_updated_at {
1505 self.vertex_updated_at.insert(*vid, *ts); }
1507
1508 for (eid, ts) in &other.edge_created_at {
1509 self.edge_created_at.entry(*eid).or_insert(*ts); }
1511 for (eid, ts) in &other.edge_updated_at {
1512 self.edge_updated_at.insert(*eid, *ts); }
1514
1515 self.estimated_size += other.estimated_size;
1518
1519 for (key, vid) in &other.constraint_index {
1521 self.constraint_index.insert(key.clone(), *vid);
1522 }
1523
1524 for (key, vid) in &other.merge_guard_index {
1527 self.merge_guard_index.insert(key.clone(), *vid);
1528 }
1529
1530 for (key, eid) in &other.edge_constraint_index {
1532 self.edge_constraint_index.insert(key.clone(), *eid);
1533 }
1534
1535 for (vid, label) in &other.pending_embeddings {
1541 self.pending_embeddings.insert(*vid, label.clone());
1542 }
1543
1544 Ok(())
1545 }
1546
1547 #[instrument(skip(self, mutations), level = "debug")]
1551 pub fn replay_mutations(&mut self, mutations: Vec<Mutation>) -> Result<()> {
1552 trace!(count = mutations.len(), "Replaying mutations");
1553 for mutation in mutations {
1554 match mutation {
1555 Mutation::InsertVertex {
1556 vid,
1557 properties,
1558 labels,
1559 } => {
1560 self.current_version += 1;
1562 let version = self.current_version;
1563
1564 self.vertex_tombstones.remove(&vid);
1565 let tracks_extid = properties.contains_key("ext_id");
1566 let entry = self.vertex_properties.entry(vid).or_default();
1567 let old_extid = if tracks_extid {
1568 Self::extid_of(entry)
1569 } else {
1570 None
1571 };
1572 Self::merge_crdt_properties(entry, properties, self.plugin_registry.as_ref());
1573 if tracks_extid {
1574 let new_extid = Self::extid_of(
1575 self.vertex_properties.get(&vid).expect("just inserted"),
1576 );
1577 self.sync_extid_index(vid, old_extid, new_extid);
1578 }
1579 self.vertex_versions.insert(vid, version);
1580 self.graph.add_vertex(vid);
1581 self.mutation_count += 1;
1582
1583 let existing = self.vertex_labels.entry(vid).or_default();
1585 Self::append_unique_labels(existing, &labels);
1586 for label in &labels {
1587 self.label_to_vids
1588 .entry(label.clone())
1589 .or_default()
1590 .insert(vid);
1591 }
1592 }
1593 Mutation::DeleteVertex { vid, labels } => {
1594 self.current_version += 1;
1595 if !labels.is_empty() {
1597 let existing = self.vertex_labels.entry(vid).or_default();
1598 Self::append_unique_labels(existing, &labels);
1599 for label in &labels {
1600 self.label_to_vids
1601 .entry(label.clone())
1602 .or_default()
1603 .insert(vid);
1604 }
1605 }
1606 self.apply_vertex_deletion(vid);
1607 }
1608 Mutation::SetVertexLabels { vid, labels } => {
1609 self.current_version += 1;
1614 self.remove_vid_from_label_index(vid);
1615 self.vertex_labels.insert(vid, labels.clone());
1616 self.index_labels_for_vid(vid, &labels);
1617 self.vertex_label_overwrites.insert(vid);
1624 self.mutation_count += 1;
1625 }
1626 Mutation::InsertEdge {
1627 src_vid,
1628 dst_vid,
1629 edge_type,
1630 eid,
1631 version: _,
1632 properties,
1633 edge_type_name,
1634 } => {
1635 self.current_version += 1;
1636 match self.apply_edge_insertion(src_vid, dst_vid, edge_type, eid, properties) {
1641 Ok(()) => {
1642 if let Some(name) = edge_type_name {
1644 self.edge_types.insert(eid, name);
1645 }
1646 }
1647 Err(e) => {
1648 tracing::warn!(
1649 ?eid,
1650 ?src_vid,
1651 ?dst_vid,
1652 error = %e,
1653 "WAL replay: skipping edge insertion to a deleted endpoint (issue #77)"
1654 );
1655 }
1656 }
1657 }
1658 Mutation::DeleteEdge {
1659 eid,
1660 src_vid,
1661 dst_vid,
1662 edge_type,
1663 version: _,
1664 } => {
1665 self.current_version += 1;
1666 self.apply_edge_deletion(eid, src_vid, dst_vid, edge_type);
1667 }
1668 }
1669 }
1670 Ok(())
1671 }
1672}
1673
1674#[cfg(test)]
1675mod tests {
1676 use super::*;
1677
1678 #[test]
1679 fn test_l0_buffer_ops() -> Result<()> {
1680 let mut l0 = L0Buffer::new(0, None);
1681 let vid_a = Vid::new(1);
1682 let vid_b = Vid::new(2);
1683 let eid_ab = Eid::new(101);
1684
1685 l0.insert_edge(vid_a, vid_b, 1, eid_ab, HashMap::new(), None)?;
1686
1687 let neighbors = l0.get_neighbors(vid_a, 1, Direction::Outgoing);
1688 assert_eq!(neighbors.len(), 1);
1689 assert_eq!(neighbors[0].0, vid_b);
1690 assert_eq!(neighbors[0].1, eid_ab);
1691
1692 l0.delete_edge(eid_ab, vid_a, vid_b, 1)?;
1693 assert!(l0.is_tombstoned(eid_ab));
1694
1695 let neighbors_after = l0.get_neighbors(vid_a, 1, Direction::Outgoing);
1697 assert_eq!(neighbors_after.len(), 0);
1698
1699 Ok(())
1700 }
1701
1702 #[test]
1708 fn validate_merge_rejects_edge_to_tombstoned_endpoint() {
1709 let mut main = L0Buffer::new(0, None);
1710 let vid_a = Vid::new(1);
1711 let vid_b = Vid::new(2);
1712 main.insert_vertex(vid_a, HashMap::new());
1713 main.insert_vertex(vid_b, HashMap::new());
1714 main.delete_vertex(vid_b).unwrap(); let mut tx = L0Buffer::new(0, None);
1718 let eid = Eid::new(101);
1719 tx.insert_edge(vid_a, vid_b, 1, eid, HashMap::new(), None)
1720 .unwrap();
1721
1722 assert!(
1723 main.validate_merge_edge_endpoints(&tx).is_err(),
1724 "edge to a tombstoned endpoint must be rejected before merge"
1725 );
1726 assert!(
1728 main.merge(&tx).is_err(),
1729 "merge must reject, not bail mid-apply"
1730 );
1731 assert!(
1732 !main.edge_endpoints.contains_key(&eid),
1733 "a rejected merge must not have partially applied the edge"
1734 );
1735 }
1736
1737 #[test]
1740 fn validate_merge_allows_edge_when_endpoint_reinserted() {
1741 let mut main = L0Buffer::new(0, None);
1742 let vid_a = Vid::new(1);
1743 let vid_b = Vid::new(2);
1744 main.insert_vertex(vid_a, HashMap::new());
1745 main.insert_vertex(vid_b, HashMap::new());
1746 main.delete_vertex(vid_b).unwrap();
1747
1748 let mut tx = L0Buffer::new(0, None);
1749 tx.insert_vertex(vid_b, HashMap::new()); let eid = Eid::new(101);
1751 tx.insert_edge(vid_a, vid_b, 1, eid, HashMap::new(), None)
1752 .unwrap();
1753
1754 assert!(main.validate_merge_edge_endpoints(&tx).is_ok());
1755 assert!(main.merge(&tx).is_ok());
1756 assert!(main.edge_endpoints.contains_key(&eid));
1757 }
1758
1759 #[test]
1761 fn validate_merge_allows_edge_to_live_endpoints() {
1762 let mut main = L0Buffer::new(0, None);
1763 let vid_a = Vid::new(1);
1764 let vid_b = Vid::new(2);
1765 main.insert_vertex(vid_a, HashMap::new());
1766 main.insert_vertex(vid_b, HashMap::new());
1767
1768 let mut tx = L0Buffer::new(0, None);
1769 let eid = Eid::new(101);
1770 tx.insert_edge(vid_a, vid_b, 1, eid, HashMap::new(), None)
1771 .unwrap();
1772
1773 assert!(main.validate_merge_edge_endpoints(&tx).is_ok());
1774 assert!(main.merge(&tx).is_ok());
1775 assert!(main.edge_endpoints.contains_key(&eid));
1776 }
1777
1778 #[test]
1779 fn test_l0_buffer_multiple_edges() -> Result<()> {
1780 let mut l0 = L0Buffer::new(0, None);
1781 let vid_a = Vid::new(1);
1782 let vid_b = Vid::new(2);
1783 let vid_c = Vid::new(3);
1784 let eid_ab = Eid::new(101);
1785 let eid_ac = Eid::new(102);
1786
1787 l0.insert_edge(vid_a, vid_b, 1, eid_ab, HashMap::new(), None)?;
1788 l0.insert_edge(vid_a, vid_c, 1, eid_ac, HashMap::new(), None)?;
1789
1790 let neighbors = l0.get_neighbors(vid_a, 1, Direction::Outgoing);
1791 assert_eq!(neighbors.len(), 2);
1792
1793 l0.delete_edge(eid_ab, vid_a, vid_b, 1)?;
1795
1796 let neighbors_after = l0.get_neighbors(vid_a, 1, Direction::Outgoing);
1798 assert_eq!(neighbors_after.len(), 1);
1799 assert_eq!(neighbors_after[0].0, vid_c);
1800
1801 Ok(())
1802 }
1803
1804 #[test]
1805 fn test_l0_buffer_edge_type_filter() -> Result<()> {
1806 let mut l0 = L0Buffer::new(0, None);
1807 let vid_a = Vid::new(1);
1808 let vid_b = Vid::new(2);
1809 let vid_c = Vid::new(3);
1810 let eid_ab = Eid::new(101);
1811 let eid_ac = Eid::new(201); l0.insert_edge(vid_a, vid_b, 1, eid_ab, HashMap::new(), None)?;
1814 l0.insert_edge(vid_a, vid_c, 2, eid_ac, HashMap::new(), None)?;
1815
1816 let type1_neighbors = l0.get_neighbors(vid_a, 1, Direction::Outgoing);
1818 assert_eq!(type1_neighbors.len(), 1);
1819 assert_eq!(type1_neighbors[0].0, vid_b);
1820
1821 let type2_neighbors = l0.get_neighbors(vid_a, 2, Direction::Outgoing);
1823 assert_eq!(type2_neighbors.len(), 1);
1824 assert_eq!(type2_neighbors[0].0, vid_c);
1825
1826 Ok(())
1827 }
1828
1829 #[test]
1830 fn test_l0_buffer_incoming_edges() -> Result<()> {
1831 let mut l0 = L0Buffer::new(0, None);
1832 let vid_a = Vid::new(1);
1833 let vid_b = Vid::new(2);
1834 let vid_c = Vid::new(3);
1835 let eid_ab = Eid::new(101);
1836 let eid_cb = Eid::new(102);
1837
1838 l0.insert_edge(vid_a, vid_b, 1, eid_ab, HashMap::new(), None)?;
1840 l0.insert_edge(vid_c, vid_b, 1, eid_cb, HashMap::new(), None)?;
1841
1842 let incoming = l0.get_neighbors(vid_b, 1, Direction::Incoming);
1844 assert_eq!(incoming.len(), 2);
1845
1846 Ok(())
1847 }
1848
1849 #[test]
1851 fn test_merge_empty_props_edge() -> Result<()> {
1852 let mut main_l0 = L0Buffer::new(0, None);
1853 let mut tx_l0 = L0Buffer::new(0, None);
1854
1855 let vid_a = Vid::new(1);
1856 let vid_b = Vid::new(2);
1857 let eid_ab = Eid::new(101);
1858
1859 tx_l0.insert_edge(vid_a, vid_b, 1, eid_ab, HashMap::new(), None)?;
1861
1862 assert!(tx_l0.edge_endpoints.contains_key(&eid_ab));
1864 assert!(!tx_l0.edge_properties.contains_key(&eid_ab)); main_l0.merge(&tx_l0)?;
1868
1869 assert!(main_l0.edge_endpoints.contains_key(&eid_ab));
1871 let neighbors = main_l0.get_neighbors(vid_a, 1, Direction::Outgoing);
1872 assert_eq!(neighbors.len(), 1);
1873 assert_eq!(neighbors[0].0, vid_b);
1874
1875 Ok(())
1876 }
1877
1878 #[test]
1880 fn test_replay_crdt_merge() -> Result<()> {
1881 use crate::runtime::wal::Mutation;
1882 use serde_json::json;
1883 use uni_common::Value;
1884
1885 let mut l0 = L0Buffer::new(0, None);
1886 let vid = Vid::new(1);
1887
1888 let counter1: Value = json!({
1891 "t": "gc",
1892 "d": {"counts": {"node1": 5}}
1893 })
1894 .into();
1895 let counter2: Value = json!({
1896 "t": "gc",
1897 "d": {"counts": {"node2": 3}}
1898 })
1899 .into();
1900
1901 let mut props1 = HashMap::new();
1903 props1.insert("counter".to_string(), counter1.clone());
1904 l0.replay_mutations(vec![Mutation::InsertVertex {
1905 vid,
1906 properties: props1,
1907 labels: vec![],
1908 }])?;
1909
1910 let mut props2 = HashMap::new();
1912 props2.insert("counter".to_string(), counter2.clone());
1913 l0.replay_mutations(vec![Mutation::InsertVertex {
1914 vid,
1915 properties: props2,
1916 labels: vec![],
1917 }])?;
1918
1919 let stored_props = l0.vertex_properties.get(&vid).unwrap();
1921 let stored_counter = stored_props.get("counter").unwrap();
1922
1923 let stored_json: serde_json::Value = stored_counter.clone().into();
1925 let data = stored_json.get("d").unwrap();
1927 let counts = data.get("counts").unwrap();
1928 assert_eq!(counts.get("node1"), Some(&json!(5)));
1929 assert_eq!(counts.get("node2"), Some(&json!(3)));
1930
1931 Ok(())
1932 }
1933
1934 #[test]
1935 fn test_merge_preserves_vertex_timestamps() -> Result<()> {
1936 let mut l0_main = L0Buffer::new(0, None);
1937 let mut l0_tx = L0Buffer::new(0, None);
1938 let vid = Vid::new(1);
1939
1940 let ts_main_created = 1000;
1942 let ts_main_updated = 1100;
1943 l0_main.insert_vertex(vid, HashMap::new());
1944 l0_main.vertex_created_at.insert(vid, ts_main_created);
1945 l0_main.vertex_updated_at.insert(vid, ts_main_updated);
1946
1947 let ts_tx_created = 2000; let ts_tx_updated = 2100; l0_tx.insert_vertex(vid, HashMap::new());
1951 l0_tx.vertex_created_at.insert(vid, ts_tx_created);
1952 l0_tx.vertex_updated_at.insert(vid, ts_tx_updated);
1953
1954 l0_main.merge(&l0_tx)?;
1956
1957 assert_eq!(
1959 *l0_main.vertex_created_at.get(&vid).unwrap(),
1960 ts_main_created,
1961 "created_at should preserve oldest timestamp"
1962 );
1963
1964 assert_eq!(
1966 *l0_main.vertex_updated_at.get(&vid).unwrap(),
1967 ts_tx_updated,
1968 "updated_at should use latest timestamp"
1969 );
1970
1971 Ok(())
1972 }
1973
1974 #[test]
1975 fn test_merge_preserves_edge_timestamps() -> Result<()> {
1976 let mut l0_main = L0Buffer::new(0, None);
1977 let mut l0_tx = L0Buffer::new(0, None);
1978 let vid_a = Vid::new(1);
1979 let vid_b = Vid::new(2);
1980 let eid = Eid::new(100);
1981
1982 let ts_main_created = 1000;
1984 let ts_main_updated = 1100;
1985 l0_main.insert_edge(vid_a, vid_b, 1, eid, HashMap::new(), None)?;
1986 l0_main.edge_created_at.insert(eid, ts_main_created);
1987 l0_main.edge_updated_at.insert(eid, ts_main_updated);
1988
1989 let ts_tx_created = 2000; let ts_tx_updated = 2100; l0_tx.insert_edge(vid_a, vid_b, 1, eid, HashMap::new(), None)?;
1993 l0_tx.edge_created_at.insert(eid, ts_tx_created);
1994 l0_tx.edge_updated_at.insert(eid, ts_tx_updated);
1995
1996 l0_main.merge(&l0_tx)?;
1998
1999 assert_eq!(
2001 *l0_main.edge_created_at.get(&eid).unwrap(),
2002 ts_main_created,
2003 "edge created_at should preserve oldest timestamp"
2004 );
2005
2006 assert_eq!(
2008 *l0_main.edge_updated_at.get(&eid).unwrap(),
2009 ts_tx_updated,
2010 "edge updated_at should use latest timestamp"
2011 );
2012
2013 Ok(())
2014 }
2015
2016 #[test]
2017 fn test_merge_created_at_not_overwritten_for_existing_vertex() -> Result<()> {
2018 use uni_common::Value;
2019
2020 let mut l0_main = L0Buffer::new(0, None);
2021 let mut l0_tx = L0Buffer::new(0, None);
2022 let vid = Vid::new(1);
2023
2024 let ts_original = 1000;
2026 l0_main.insert_vertex(vid, HashMap::new());
2027 l0_main.vertex_created_at.insert(vid, ts_original);
2028 l0_main.vertex_updated_at.insert(vid, ts_original);
2029
2030 let ts_tx = 2000;
2032 let mut props = HashMap::new();
2033 props.insert("updated".to_string(), Value::String("yes".to_string()));
2034 l0_tx.insert_vertex(vid, props);
2035 l0_tx.vertex_created_at.insert(vid, ts_tx);
2036 l0_tx.vertex_updated_at.insert(vid, ts_tx);
2037
2038 l0_main.merge(&l0_tx)?;
2040
2041 assert_eq!(
2043 *l0_main.vertex_created_at.get(&vid).unwrap(),
2044 ts_original,
2045 "created_at must not be overwritten for existing vertex"
2046 );
2047
2048 assert_eq!(
2050 *l0_main.vertex_updated_at.get(&vid).unwrap(),
2051 ts_tx,
2052 "updated_at should reflect transaction timestamp"
2053 );
2054
2055 assert!(
2057 l0_main
2058 .vertex_properties
2059 .get(&vid)
2060 .unwrap()
2061 .contains_key("updated")
2062 );
2063
2064 Ok(())
2065 }
2066
2067 #[test]
2069 fn test_replay_mutations_preserves_vertex_labels() -> Result<()> {
2070 use crate::runtime::wal::Mutation;
2071
2072 let mut l0 = L0Buffer::new(0, None);
2073 let vid = Vid::new(42);
2074
2075 let mutations = vec![Mutation::InsertVertex {
2077 vid,
2078 properties: {
2079 let mut props = HashMap::new();
2080 props.insert(
2081 "name".to_string(),
2082 uni_common::Value::String("Alice".to_string()),
2083 );
2084 props
2085 },
2086 labels: vec!["Person".to_string(), "User".to_string()],
2087 }];
2088
2089 l0.replay_mutations(mutations)?;
2091
2092 assert!(l0.vertex_properties.contains_key(&vid));
2094
2095 let labels = l0.get_vertex_labels(vid).expect("Labels should exist");
2097 assert_eq!(labels.len(), 2);
2098 assert!(labels.contains(&"Person".to_string()));
2099 assert!(labels.contains(&"User".to_string()));
2100
2101 let person_vids = l0.vids_for_label("Person");
2103 assert_eq!(person_vids.len(), 1);
2104 assert_eq!(person_vids[0], vid);
2105
2106 let user_vids = l0.vids_for_label("User");
2107 assert_eq!(user_vids.len(), 1);
2108 assert_eq!(user_vids[0], vid);
2109
2110 Ok(())
2111 }
2112
2113 #[test]
2115 fn test_replay_mutations_preserves_delete_vertex_labels() -> Result<()> {
2116 use crate::runtime::wal::Mutation;
2117
2118 let mut l0 = L0Buffer::new(0, None);
2119 let vid = Vid::new(99);
2120
2121 l0.insert_vertex_with_labels(
2123 vid,
2124 HashMap::new(),
2125 &["Person".to_string(), "Admin".to_string()],
2126 );
2127
2128 assert!(l0.vertex_properties.contains_key(&vid));
2130 let labels = l0.get_vertex_labels(vid).expect("Labels should exist");
2131 assert_eq!(labels.len(), 2);
2132
2133 let mutations = vec![Mutation::DeleteVertex {
2135 vid,
2136 labels: vec!["Person".to_string(), "Admin".to_string()],
2137 }];
2138
2139 l0.replay_mutations(mutations)?;
2141
2142 assert!(l0.vertex_tombstones.contains(&vid));
2144
2145 let labels = l0.get_vertex_labels(vid);
2148 assert!(
2149 labels.is_some(),
2150 "Labels should be preserved even after deletion for tombstone flushing"
2151 );
2152
2153 Ok(())
2154 }
2155
2156 #[test]
2158 fn test_replay_mutations_preserves_edge_type_name() -> Result<()> {
2159 use crate::runtime::wal::Mutation;
2160
2161 let mut l0 = L0Buffer::new(0, None);
2162 let src = Vid::new(1);
2163 let dst = Vid::new(2);
2164 let eid = Eid::new(500);
2165 let edge_type = 100;
2166
2167 let mutations = vec![Mutation::InsertEdge {
2169 src_vid: src,
2170 dst_vid: dst,
2171 edge_type,
2172 eid,
2173 version: 1,
2174 properties: {
2175 let mut props = HashMap::new();
2176 props.insert("since".to_string(), uni_common::Value::Int(2020));
2177 props
2178 },
2179 edge_type_name: Some("KNOWS".to_string()),
2180 }];
2181
2182 l0.replay_mutations(mutations)?;
2184
2185 assert!(l0.edge_endpoints.contains_key(&eid));
2187
2188 let type_name = l0.get_edge_type(eid).expect("Edge type name should exist");
2190 assert_eq!(type_name, "KNOWS");
2191
2192 let knows_eids = l0.eids_for_type("KNOWS");
2194 assert_eq!(knows_eids.len(), 1);
2195 assert_eq!(knows_eids[0], eid);
2196
2197 Ok(())
2198 }
2199
2200 #[test]
2202 fn test_edge_type_mapping_survives_multiple_replays() -> Result<()> {
2203 use crate::runtime::wal::Mutation;
2204
2205 let mut l0 = L0Buffer::new(0, None);
2206
2207 let mutations = vec![
2209 Mutation::InsertEdge {
2210 src_vid: Vid::new(1),
2211 dst_vid: Vid::new(2),
2212 edge_type: 100,
2213 eid: Eid::new(1000),
2214 version: 1,
2215 properties: HashMap::new(),
2216 edge_type_name: Some("KNOWS".to_string()),
2217 },
2218 Mutation::InsertEdge {
2219 src_vid: Vid::new(2),
2220 dst_vid: Vid::new(3),
2221 edge_type: 101,
2222 eid: Eid::new(1001),
2223 version: 2,
2224 properties: HashMap::new(),
2225 edge_type_name: Some("LIKES".to_string()),
2226 },
2227 Mutation::InsertEdge {
2228 src_vid: Vid::new(3),
2229 dst_vid: Vid::new(1),
2230 edge_type: 100,
2231 eid: Eid::new(1002),
2232 version: 3,
2233 properties: HashMap::new(),
2234 edge_type_name: Some("KNOWS".to_string()),
2235 },
2236 ];
2237
2238 l0.replay_mutations(mutations)?;
2239
2240 assert_eq!(l0.get_edge_type(Eid::new(1000)), Some("KNOWS"));
2242 assert_eq!(l0.get_edge_type(Eid::new(1001)), Some("LIKES"));
2243 assert_eq!(l0.get_edge_type(Eid::new(1002)), Some("KNOWS"));
2244
2245 let knows_edges = l0.eids_for_type("KNOWS");
2247 assert_eq!(knows_edges.len(), 2);
2248 assert!(knows_edges.contains(&Eid::new(1000)));
2249 assert!(knows_edges.contains(&Eid::new(1002)));
2250
2251 let likes_edges = l0.eids_for_type("LIKES");
2252 assert_eq!(likes_edges.len(), 1);
2253 assert_eq!(likes_edges[0], Eid::new(1001));
2254
2255 Ok(())
2256 }
2257
2258 #[test]
2260 fn test_replay_mutations_combined_labels_and_edge_types() -> Result<()> {
2261 use crate::runtime::wal::Mutation;
2262
2263 let mut l0 = L0Buffer::new(0, None);
2264 let alice = Vid::new(1);
2265 let bob = Vid::new(2);
2266 let eid = Eid::new(100);
2267
2268 let mutations = vec![
2270 Mutation::InsertVertex {
2272 vid: alice,
2273 properties: {
2274 let mut props = HashMap::new();
2275 props.insert(
2276 "name".to_string(),
2277 uni_common::Value::String("Alice".to_string()),
2278 );
2279 props
2280 },
2281 labels: vec!["Person".to_string()],
2282 },
2283 Mutation::InsertVertex {
2285 vid: bob,
2286 properties: {
2287 let mut props = HashMap::new();
2288 props.insert(
2289 "name".to_string(),
2290 uni_common::Value::String("Bob".to_string()),
2291 );
2292 props
2293 },
2294 labels: vec!["Person".to_string()],
2295 },
2296 Mutation::InsertEdge {
2298 src_vid: alice,
2299 dst_vid: bob,
2300 edge_type: 1,
2301 eid,
2302 version: 3,
2303 properties: HashMap::new(),
2304 edge_type_name: Some("KNOWS".to_string()),
2305 },
2306 ];
2307
2308 l0.replay_mutations(mutations)?;
2310
2311 assert_eq!(l0.get_vertex_labels(alice).unwrap().len(), 1);
2313 assert_eq!(l0.get_vertex_labels(bob).unwrap().len(), 1);
2314 assert_eq!(l0.vids_for_label("Person").len(), 2);
2315
2316 assert_eq!(l0.get_edge_type(eid).unwrap(), "KNOWS");
2318 assert_eq!(l0.eids_for_type("KNOWS").len(), 1);
2319
2320 let alice_neighbors = l0.get_neighbors(alice, 1, Direction::Outgoing);
2322 assert_eq!(alice_neighbors.len(), 1);
2323 assert_eq!(alice_neighbors[0].0, bob);
2324
2325 Ok(())
2326 }
2327
2328 #[test]
2330 fn test_replay_mutations_backward_compat_empty_labels() -> Result<()> {
2331 use crate::runtime::wal::Mutation;
2332
2333 let mut l0 = L0Buffer::new(0, None);
2334 let vid = Vid::new(1);
2335
2336 let mutations = vec![Mutation::InsertVertex {
2339 vid,
2340 properties: HashMap::new(),
2341 labels: vec![], }];
2343
2344 l0.replay_mutations(mutations)?;
2345
2346 assert!(l0.vertex_properties.contains_key(&vid));
2348
2349 let labels = l0.get_vertex_labels(vid);
2351 assert!(labels.is_some(), "Labels entry should exist even if empty");
2352 assert_eq!(labels.unwrap().len(), 0);
2353
2354 Ok(())
2355 }
2356
2357 #[test]
2358 fn test_now_nanos_returns_nanosecond_range() {
2359 let now = now_nanos();
2363
2364 assert!(
2366 now > 1_700_000_000_000_000_000,
2367 "now_nanos() returned {}, expected > 1.7e18 for nanoseconds",
2368 now
2369 );
2370
2371 assert!(
2373 now < 4_100_000_000_000_000_000,
2374 "now_nanos() returned {}, expected < 4.1e18",
2375 now
2376 );
2377 }
2378}