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_tombstones.remove(&vid);
610
611 self.vertex_partial_keys.remove(&vid);
614
615 let props_size = Self::estimate_properties_size(&properties);
618 let props_count = properties.len();
619 let tracks_extid = properties.contains_key("ext_id");
620
621 let entry = self.vertex_properties.entry(vid).or_default();
622 let old_extid = if tracks_extid {
623 Self::extid_of(entry)
624 } else {
625 None
626 };
627 Self::merge_crdt_properties(entry, properties, self.plugin_registry.as_ref());
628 if tracks_extid {
629 let new_extid =
630 Self::extid_of(self.vertex_properties.get(&vid).expect("just inserted"));
631 self.sync_extid_index(vid, old_extid, new_extid);
632 }
633 self.vertex_versions.insert(vid, version);
634
635 self.vertex_created_at.entry(vid).or_insert(now);
637 self.vertex_updated_at.insert(vid, now);
638
639 let labels_size: usize = labels.iter().map(|l| l.len() + 24).sum();
642 let existing = self.vertex_labels.entry(vid).or_default();
643 Self::append_unique_labels(existing, labels);
644 self.index_labels_for_vid(vid, labels);
645
646 self.graph.add_vertex(vid);
647 self.mutation_count += 1;
648 self.mutation_stats.nodes_created += 1;
649 self.mutation_stats.properties_set += props_count;
650 self.mutation_stats.labels_added += labels.len();
651
652 self.estimated_size += 8 + props_size + 16 + labels_size + 32;
653 }
654
655 pub fn insert_vertex_partial_full(
676 &mut self,
677 vid: Vid,
678 props: Properties,
679 touched_keys: HashSet<String>,
680 labels: &[String],
681 ) {
682 self.insert_vertex_with_labels_partial_impl(vid, props, labels, false);
686 self.vertex_partial_keys
687 .entry(vid)
688 .or_default()
689 .extend(touched_keys);
690 }
691
692 pub fn insert_vertex_partial(&mut self, vid: Vid, touched: Properties, labels: &[String]) {
696 let touched_keys: Vec<String> = touched.keys().cloned().collect();
700
701 let already_full = self.vertex_properties.contains_key(&vid)
712 && !self.vertex_partial_keys.contains_key(&vid);
713
714 self.insert_vertex_with_labels_partial_impl(vid, touched, labels, false);
719
720 if !already_full {
721 self.vertex_partial_keys
722 .entry(vid)
723 .or_default()
724 .extend(touched_keys);
725 }
726 }
727
728 fn insert_vertex_with_labels_partial_impl(
732 &mut self,
733 vid: Vid,
734 properties: Properties,
735 labels: &[String],
736 skip_wal: bool,
737 ) {
738 self.current_version += 1;
739 let version = self.current_version;
740 let now = now_nanos();
741
742 if !skip_wal && let Some(wal) = &self.wal {
743 let _ = wal.append(Mutation::InsertVertex {
749 vid,
750 properties: properties.clone(),
751 labels: labels.to_vec(),
752 });
753 }
754
755 self.vertex_tombstones.remove(&vid);
756 let props_size = Self::estimate_properties_size(&properties);
762 let props_count = properties.len();
763 let tracks_extid = properties.contains_key("ext_id");
764
765 let entry = self.vertex_properties.entry(vid).or_default();
766 let old_extid = if tracks_extid {
767 Self::extid_of(entry)
768 } else {
769 None
770 };
771 Self::merge_crdt_properties(entry, properties, self.plugin_registry.as_ref());
772 if tracks_extid {
773 let new_extid =
774 Self::extid_of(self.vertex_properties.get(&vid).expect("just inserted"));
775 self.sync_extid_index(vid, old_extid, new_extid);
776 }
777 self.vertex_versions.insert(vid, version);
778
779 self.vertex_created_at.entry(vid).or_insert(now);
780 self.vertex_updated_at.insert(vid, now);
781
782 let labels_size: usize = labels.iter().map(|l| l.len() + 24).sum();
783 let existing = self.vertex_labels.entry(vid).or_default();
784 Self::append_unique_labels(existing, labels);
785 self.index_labels_for_vid(vid, labels);
786
787 self.graph.add_vertex(vid);
788 self.mutation_count += 1;
789 self.mutation_stats.properties_set += props_count;
792 self.mutation_stats.labels_added += labels.len();
793
794 self.estimated_size += 8 + props_size + 16 + labels_size + 32;
795 }
796
797 pub fn add_vertex_labels(&mut self, vid: Vid, labels: &[String]) {
799 let existing = self.vertex_labels.entry(vid).or_default();
800 Self::append_unique_labels(existing, labels);
801 self.index_labels_for_vid(vid, labels);
802 }
803
804 pub fn remove_vertex_label(&mut self, vid: Vid, label: &str) -> bool {
807 if let Some(labels) = self.vertex_labels.get_mut(&vid)
808 && let Some(pos) = labels.iter().position(|l| l == label)
809 {
810 labels.remove(pos);
811 if let Some(set) = self.label_to_vids.get_mut(label) {
812 set.remove(&vid);
813 }
814 self.current_version += 1;
815 self.mutation_count += 1;
816 self.mutation_stats.labels_removed += 1;
817 return true;
820 }
821 false
822 }
823
824 pub fn set_edge_type(&mut self, eid: Eid, edge_type: String) {
826 self.edge_types.insert(eid, edge_type);
827 }
828
829 pub fn delete_vertex(&mut self, vid: Vid) -> Result<()> {
830 self.delete_vertex_impl(vid, false)
831 }
832
833 fn delete_vertex_impl(&mut self, vid: Vid, skip_wal: bool) -> Result<()> {
836 self.current_version += 1;
837
838 if !skip_wal && let Some(wal) = &mut self.wal {
839 let labels = self.vertex_labels.get(&vid).cloned().unwrap_or_default();
840 wal.append(Mutation::DeleteVertex { vid, labels })?;
841 }
842
843 self.apply_vertex_deletion(vid);
844 Ok(())
845 }
846
847 fn apply_vertex_deletion(&mut self, vid: Vid) {
851 let version = self.current_version;
852
853 let mut edges_to_remove = HashSet::new();
855
856 for entry in self.graph.neighbors(vid, Direction::Outgoing) {
858 edges_to_remove.insert(entry.eid);
859 }
860
861 for entry in self.graph.neighbors(vid, Direction::Incoming) {
863 edges_to_remove.insert(entry.eid); }
865
866 let cascaded_edges_count = edges_to_remove.len();
867
868 for eid in edges_to_remove {
870 if let Some((src, dst, etype)) = self.edge_endpoints.get(&eid) {
872 self.tombstones.insert(
873 eid,
874 TombstoneEntry {
875 eid,
876 src_vid: *src,
877 dst_vid: *dst,
878 edge_type: *etype,
879 },
880 );
881 self.edge_versions.insert(eid, version);
882 self.edge_endpoints.remove(&eid);
883 self.edge_properties.remove(&eid);
884 self.graph.remove_edge(eid);
885 self.mutation_count += 1;
886 self.mutation_stats.relationships_deleted += 1;
887 }
888 }
889
890 self.remove_vid_from_label_index(vid);
891 self.vertex_tombstones.insert(vid);
892 if let Some(props) = self.vertex_properties.get(&vid)
895 && let Some(ext) = Self::extid_of(props)
896 && self.extid_index.get(&ext) == Some(&vid)
897 {
898 self.extid_index.remove(&ext);
899 }
900 self.vertex_properties.remove(&vid);
901 self.vertex_partial_keys.remove(&vid);
903 self.vertex_versions.insert(vid, version);
904 self.graph.remove_vertex(vid);
905 self.mutation_count += 1;
906 self.mutation_stats.nodes_deleted += 1;
907
908 self.constraint_index.retain(|_, v| *v != vid);
910 self.merge_guard_index.retain(|_, v| *v != vid);
913
914 self.estimated_size += cascaded_edges_count * 72 + 8;
916 }
917
918 pub fn insert_edge(
919 &mut self,
920 src_vid: Vid,
921 dst_vid: Vid,
922 edge_type: u32,
923 eid: Eid,
924 properties: Properties,
925 edge_type_name: Option<String>,
926 ) -> Result<()> {
927 self.insert_edge_impl(
928 src_vid,
929 dst_vid,
930 edge_type,
931 eid,
932 properties,
933 edge_type_name,
934 false,
935 )
936 }
937
938 #[allow(clippy::too_many_arguments)]
941 fn insert_edge_impl(
942 &mut self,
943 src_vid: Vid,
944 dst_vid: Vid,
945 edge_type: u32,
946 eid: Eid,
947 properties: Properties,
948 edge_type_name: Option<String>,
949 skip_wal: bool,
950 ) -> Result<()> {
951 self.current_version += 1;
952 let now = now_nanos();
953
954 if !skip_wal && let Some(wal) = &mut self.wal {
955 wal.append(Mutation::InsertEdge {
956 src_vid,
957 dst_vid,
958 edge_type,
959 eid,
960 version: self.current_version,
961 properties: properties.clone(),
962 edge_type_name: edge_type_name.clone(),
963 })?;
964 }
965
966 self.apply_edge_insertion(src_vid, dst_vid, edge_type, eid, properties)?;
967
968 let type_name_size = if let Some(ref name) = edge_type_name {
970 let size = name.len() + 24;
971 self.edge_types.insert(eid, name.clone());
972 size
973 } else {
974 0
975 };
976
977 self.edge_created_at.entry(eid).or_insert(now);
979 self.edge_updated_at.insert(eid, now);
980
981 self.edge_partial_keys.remove(&eid);
984
985 self.estimated_size += type_name_size;
986
987 Ok(())
988 }
989
990 #[allow(clippy::too_many_arguments)]
995 pub fn insert_edge_partial_full(
996 &mut self,
997 src_vid: Vid,
998 dst_vid: Vid,
999 edge_type: u32,
1000 eid: Eid,
1001 properties: Properties,
1002 edge_type_name: Option<String>,
1003 touched_keys: HashSet<String>,
1004 ) -> Result<()> {
1005 self.current_version += 1;
1006 let now = now_nanos();
1007
1008 if let Some(wal) = &mut self.wal {
1009 wal.append(Mutation::InsertEdge {
1010 src_vid,
1011 dst_vid,
1012 edge_type,
1013 eid,
1014 version: self.current_version,
1015 properties: properties.clone(),
1016 edge_type_name: edge_type_name.clone(),
1017 })?;
1018 }
1019
1020 self.apply_edge_insertion(src_vid, dst_vid, edge_type, eid, properties)?;
1021
1022 self.edge_partial_keys
1026 .entry(eid)
1027 .or_default()
1028 .extend(touched_keys);
1029
1030 let type_name_size = if let Some(ref name) = edge_type_name {
1031 let size = name.len() + 24;
1032 self.edge_types.insert(eid, name.clone());
1033 size
1034 } else {
1035 0
1036 };
1037
1038 self.edge_created_at.entry(eid).or_insert(now);
1039 self.edge_updated_at.insert(eid, now);
1040
1041 self.estimated_size += type_name_size;
1042
1043 Ok(())
1044 }
1045
1046 fn apply_edge_insertion(
1055 &mut self,
1056 src_vid: Vid,
1057 dst_vid: Vid,
1058 edge_type: u32,
1059 eid: Eid,
1060 properties: Properties,
1061 ) -> Result<()> {
1062 let version = self.current_version;
1063
1064 if self.vertex_tombstones.contains(&src_vid) {
1067 anyhow::bail!(
1068 "Cannot insert edge: source vertex {} has been deleted (issue #77)",
1069 src_vid
1070 );
1071 }
1072 if self.vertex_tombstones.contains(&dst_vid) {
1073 anyhow::bail!(
1074 "Cannot insert edge: destination vertex {} has been deleted (issue #77)",
1075 dst_vid
1076 );
1077 }
1078
1079 if !self.graph.contains_vertex(src_vid) {
1084 self.graph.add_vertex(src_vid);
1085 }
1086 if !self.graph.contains_vertex(dst_vid) {
1087 self.graph.add_vertex(dst_vid);
1088 }
1089
1090 self.graph.add_edge(src_vid, dst_vid, eid, edge_type);
1091
1092 let props_size = Self::estimate_properties_size(&properties);
1094 let props_count = properties.len();
1095 if !properties.is_empty() {
1096 let entry = self.edge_properties.entry(eid).or_default();
1097 Self::merge_crdt_properties(entry, properties, self.plugin_registry.as_ref());
1098 }
1099
1100 self.edge_versions.insert(eid, version);
1101 self.edge_endpoints
1102 .insert(eid, (src_vid, dst_vid, edge_type));
1103 self.tombstones.remove(&eid);
1104 self.mutation_count += 1;
1105 self.mutation_stats.relationships_created += 1;
1106 self.mutation_stats.properties_set += props_count;
1107
1108 self.estimated_size += 24 + props_size + 16 + 28 + 32;
1110
1111 Ok(())
1112 }
1113
1114 pub fn delete_edge(
1115 &mut self,
1116 eid: Eid,
1117 src_vid: Vid,
1118 dst_vid: Vid,
1119 edge_type: u32,
1120 ) -> Result<()> {
1121 self.delete_edge_impl(eid, src_vid, dst_vid, edge_type, false)
1122 }
1123
1124 fn delete_edge_impl(
1127 &mut self,
1128 eid: Eid,
1129 src_vid: Vid,
1130 dst_vid: Vid,
1131 edge_type: u32,
1132 skip_wal: bool,
1133 ) -> Result<()> {
1134 self.current_version += 1;
1135 let now = now_nanos();
1136
1137 if !skip_wal && let Some(wal) = &mut self.wal {
1138 wal.append(Mutation::DeleteEdge {
1139 eid,
1140 src_vid,
1141 dst_vid,
1142 edge_type,
1143 version: self.current_version,
1144 })?;
1145 }
1146
1147 self.apply_edge_deletion(eid, src_vid, dst_vid, edge_type);
1148
1149 self.edge_updated_at.insert(eid, now);
1151
1152 Ok(())
1153 }
1154
1155 fn apply_edge_deletion(&mut self, eid: Eid, src_vid: Vid, dst_vid: Vid, edge_type: u32) {
1159 let version = self.current_version;
1160
1161 self.tombstones.insert(
1162 eid,
1163 TombstoneEntry {
1164 eid,
1165 src_vid,
1166 dst_vid,
1167 edge_type,
1168 },
1169 );
1170 self.edge_versions.insert(eid, version);
1171 self.edge_partial_keys.remove(&eid);
1174 self.edge_constraint_index.retain(|_, e| *e != eid);
1177 self.graph.remove_edge(eid);
1178 self.mutation_count += 1;
1179 self.mutation_stats.relationships_deleted += 1;
1180
1181 self.estimated_size += 80;
1183 }
1184
1185 pub fn get_neighbors(
1188 &self,
1189 vid: Vid,
1190 edge_type: u32,
1191 direction: Direction,
1192 ) -> Vec<(Vid, Eid, u64)> {
1193 let edges = self.graph.neighbors(vid, direction);
1194
1195 edges
1196 .iter()
1197 .filter(|e| e.edge_type == edge_type && !self.is_tombstoned(e.eid))
1198 .map(|e| {
1199 let neighbor = match direction {
1200 Direction::Outgoing => e.dst_vid,
1201 Direction::Incoming => e.src_vid,
1202 };
1203 let version = self.edge_versions.get(&e.eid).copied().unwrap_or(0);
1204 (neighbor, e.eid, version)
1205 })
1206 .collect()
1207 }
1208
1209 pub fn is_tombstoned(&self, eid: Eid) -> bool {
1210 self.tombstones.contains_key(&eid)
1211 }
1212
1213 pub fn vids_for_label(&self, label_name: &str) -> Vec<Vid> {
1216 self.label_to_vids
1217 .get(label_name)
1218 .map(|set| set.iter().copied().collect())
1219 .unwrap_or_default()
1220 }
1221
1222 pub fn all_vertex_vids(&self) -> Vec<Vid> {
1226 self.vertex_properties.keys().copied().collect()
1227 }
1228
1229 pub fn vids_for_labels(&self, label_names: &[&str]) -> Vec<Vid> {
1232 let mut result = HashSet::new();
1233 for label_name in label_names {
1234 if let Some(set) = self.label_to_vids.get(*label_name) {
1235 result.extend(set.iter().copied());
1236 }
1237 }
1238 result.into_iter().collect()
1239 }
1240
1241 pub fn vids_with_all_labels(&self, label_names: &[&str]) -> Vec<Vid> {
1244 if label_names.is_empty() {
1245 return Vec::new();
1246 }
1247 let sets: Vec<&HashSet<Vid>> = match label_names
1250 .iter()
1251 .map(|ln| self.label_to_vids.get(*ln))
1252 .collect::<Option<Vec<_>>>()
1253 {
1254 Some(s) => s,
1255 None => return Vec::new(),
1256 };
1257 let smallest = sets.iter().min_by_key(|s| s.len()).unwrap();
1259 smallest
1260 .iter()
1261 .copied()
1262 .filter(|vid| sets.iter().all(|s| s.contains(vid)))
1263 .collect()
1264 }
1265
1266 pub fn get_vertex_labels(&self, vid: Vid) -> Option<&[String]> {
1268 self.vertex_labels.get(&vid).map(|v| v.as_slice())
1269 }
1270
1271 pub fn get_edge_type(&self, eid: Eid) -> Option<&str> {
1273 self.edge_types.get(&eid).map(|s| s.as_str())
1274 }
1275
1276 pub fn eids_for_type(&self, type_name: &str) -> Vec<Eid> {
1279 self.edge_types
1280 .iter()
1281 .filter(|(eid, etype)| *etype == type_name && !self.tombstones.contains_key(eid))
1282 .map(|(eid, _)| *eid)
1283 .collect()
1284 }
1285
1286 pub fn all_edge_eids(&self) -> Vec<Eid> {
1290 self.edge_endpoints
1291 .keys()
1292 .filter(|eid| !self.tombstones.contains_key(eid))
1293 .copied()
1294 .collect()
1295 }
1296
1297 pub fn get_edge_endpoints(&self, eid: Eid) -> Option<(Vid, Vid)> {
1299 self.edge_endpoints
1300 .get(&eid)
1301 .map(|(src, dst, _)| (*src, *dst))
1302 }
1303
1304 pub fn get_edge_endpoint_full(&self, eid: Eid) -> Option<(Vid, Vid, u32)> {
1306 self.edge_endpoints.get(&eid).copied()
1307 }
1308
1309 pub fn insert_constraint_key(&mut self, key: Vec<u8>, vid: Vid) {
1311 self.constraint_index.insert(key, vid);
1312 }
1313
1314 pub fn has_constraint_key(&self, key: &[u8], exclude_vid: Vid) -> bool {
1317 self.constraint_index
1318 .get(key)
1319 .is_some_and(|&v| v != exclude_vid)
1320 }
1321
1322 pub fn insert_edge_constraint_key(&mut self, key: Vec<u8>, eid: Eid) {
1325 self.edge_constraint_index.insert(key, eid);
1326 }
1327
1328 pub fn has_edge_constraint_key(&self, key: &[u8], exclude_eid: Eid) -> bool {
1332 self.edge_constraint_index
1333 .get(key)
1334 .is_some_and(|&e| e != exclude_eid)
1335 }
1336
1337 pub fn insert_merge_guard_key(&mut self, key: Vec<u8>, vid: Vid) {
1339 self.merge_guard_index.insert(key, vid);
1340 }
1341
1342 pub fn has_merge_guard_key(&self, key: &[u8], exclude_vid: Vid) -> bool {
1345 self.merge_guard_index
1346 .get(key)
1347 .is_some_and(|&v| v != exclude_vid)
1348 }
1349
1350 #[instrument(skip(self, other), level = "trace")]
1351 pub fn validate_merge_edge_endpoints(&self, other: &L0Buffer) -> Result<()> {
1370 let is_deleted = |vid: &Vid| {
1374 (self.vertex_tombstones.contains(vid) || other.vertex_tombstones.contains(vid))
1375 && !other.vertex_properties.contains_key(vid)
1376 };
1377 for (eid, (src_vid, dst_vid, _etype)) in &other.edge_endpoints {
1378 if other.tombstones.contains_key(eid) {
1379 continue; }
1381 if is_deleted(src_vid) {
1382 anyhow::bail!(
1383 "Cannot insert edge {}: source vertex {} has been deleted (issue #77)",
1384 eid,
1385 src_vid
1386 );
1387 }
1388 if is_deleted(dst_vid) {
1389 anyhow::bail!(
1390 "Cannot insert edge {}: destination vertex {} has been deleted (issue #77)",
1391 eid,
1392 dst_vid
1393 );
1394 }
1395 }
1396 Ok(())
1397 }
1398
1399 pub fn merge(&mut self, other: &L0Buffer) -> Result<()> {
1400 self.validate_merge_edge_endpoints(other)?;
1404 self.merge_validated(
1405 other,
1406 other.vertex_properties.clone(),
1407 other.edge_properties.clone(),
1408 )
1409 }
1410
1411 pub fn merge_take(&mut self, other: &mut L0Buffer) -> Result<()> {
1420 self.validate_merge_edge_endpoints(other)?;
1423 let vertex_props = std::mem::take(&mut other.vertex_properties);
1424 let edge_props = std::mem::take(&mut other.edge_properties);
1425 self.merge_validated(other, vertex_props, edge_props)
1426 }
1427
1428 fn merge_validated(
1432 &mut self,
1433 other: &L0Buffer,
1434 vertex_props: HashMap<Vid, Properties>,
1435 mut edge_props: HashMap<Eid, Properties>,
1436 ) -> Result<()> {
1437 trace!(
1438 other_mutation_count = other.mutation_count,
1439 "Merging L0 buffer"
1440 );
1441 for &vid in &other.vertex_tombstones {
1446 self.delete_vertex_impl(vid, true)?;
1447 }
1448
1449 for (vid, props) in vertex_props {
1450 let labels = other.vertex_labels.get(&vid).cloned().unwrap_or_default();
1451 self.insert_vertex_with_labels_impl(vid, props, &labels, true);
1452 }
1453
1454 for (vid, labels) in &other.vertex_labels {
1456 if !self.vertex_labels.contains_key(vid) {
1457 self.vertex_labels.insert(*vid, labels.clone());
1458 for label in labels {
1459 self.label_to_vids
1460 .entry(label.clone())
1461 .or_default()
1462 .insert(*vid);
1463 }
1464 }
1465 }
1466
1467 for vid in &other.vertex_label_overwrites {
1474 if other.vertex_tombstones.contains(vid) {
1475 continue;
1476 }
1477 let labels = other.vertex_labels.get(vid).cloned().unwrap_or_default();
1478 self.remove_vid_from_label_index(*vid);
1479 self.vertex_labels.insert(*vid, labels.clone());
1480 self.index_labels_for_vid(*vid, &labels);
1481 self.vertex_label_overwrites.insert(*vid);
1486 }
1487
1488 for (eid, (src, dst, etype)) in &other.edge_endpoints {
1490 if other.tombstones.contains_key(eid) {
1491 self.delete_edge_impl(*eid, *src, *dst, *etype, true)?;
1492 } else {
1493 let props = edge_props.remove(eid).unwrap_or_default();
1494 let etype_name = other.edge_types.get(eid).cloned();
1495 self.insert_edge_impl(*src, *dst, *etype, *eid, props, etype_name, true)?;
1496 }
1497 }
1498
1499 for (eid, tombstone) in &other.tombstones {
1503 if !other.edge_endpoints.contains_key(eid) {
1504 self.delete_edge_impl(
1505 *eid,
1506 tombstone.src_vid,
1507 tombstone.dst_vid,
1508 tombstone.edge_type,
1509 true,
1510 )?;
1511 }
1512 }
1513
1514 for (vid, ts) in &other.vertex_created_at {
1519 self.vertex_created_at.entry(*vid).or_insert(*ts); }
1521 for (vid, ts) in &other.vertex_updated_at {
1522 self.vertex_updated_at.insert(*vid, *ts); }
1524
1525 for (eid, ts) in &other.edge_created_at {
1526 self.edge_created_at.entry(*eid).or_insert(*ts); }
1528 for (eid, ts) in &other.edge_updated_at {
1529 self.edge_updated_at.insert(*eid, *ts); }
1531
1532 self.estimated_size += other.estimated_size;
1535
1536 for (key, vid) in &other.constraint_index {
1538 self.constraint_index.insert(key.clone(), *vid);
1539 }
1540
1541 for (key, vid) in &other.merge_guard_index {
1544 self.merge_guard_index.insert(key.clone(), *vid);
1545 }
1546
1547 for (key, eid) in &other.edge_constraint_index {
1549 self.edge_constraint_index.insert(key.clone(), *eid);
1550 }
1551
1552 for (vid, label) in &other.pending_embeddings {
1558 self.pending_embeddings.insert(*vid, label.clone());
1559 }
1560
1561 Ok(())
1562 }
1563
1564 #[instrument(skip(self, mutations), level = "debug")]
1568 pub fn replay_mutations(&mut self, mutations: Vec<Mutation>) -> Result<()> {
1569 trace!(count = mutations.len(), "Replaying mutations");
1570 for mutation in mutations {
1571 match mutation {
1572 Mutation::InsertVertex {
1573 vid,
1574 properties,
1575 labels,
1576 } => {
1577 self.current_version += 1;
1579 let version = self.current_version;
1580
1581 self.vertex_tombstones.remove(&vid);
1582 let tracks_extid = properties.contains_key("ext_id");
1583 let entry = self.vertex_properties.entry(vid).or_default();
1584 let old_extid = if tracks_extid {
1585 Self::extid_of(entry)
1586 } else {
1587 None
1588 };
1589 Self::merge_crdt_properties(entry, properties, self.plugin_registry.as_ref());
1590 if tracks_extid {
1591 let new_extid = Self::extid_of(
1592 self.vertex_properties.get(&vid).expect("just inserted"),
1593 );
1594 self.sync_extid_index(vid, old_extid, new_extid);
1595 }
1596 self.vertex_versions.insert(vid, version);
1597 self.graph.add_vertex(vid);
1598 self.mutation_count += 1;
1599
1600 let existing = self.vertex_labels.entry(vid).or_default();
1602 Self::append_unique_labels(existing, &labels);
1603 for label in &labels {
1604 self.label_to_vids
1605 .entry(label.clone())
1606 .or_default()
1607 .insert(vid);
1608 }
1609 }
1610 Mutation::DeleteVertex { vid, labels } => {
1611 self.current_version += 1;
1612 if !labels.is_empty() {
1614 let existing = self.vertex_labels.entry(vid).or_default();
1615 Self::append_unique_labels(existing, &labels);
1616 for label in &labels {
1617 self.label_to_vids
1618 .entry(label.clone())
1619 .or_default()
1620 .insert(vid);
1621 }
1622 }
1623 self.apply_vertex_deletion(vid);
1624 }
1625 Mutation::SetVertexLabels { vid, labels } => {
1626 self.current_version += 1;
1631 self.remove_vid_from_label_index(vid);
1632 self.vertex_labels.insert(vid, labels.clone());
1633 self.index_labels_for_vid(vid, &labels);
1634 self.vertex_label_overwrites.insert(vid);
1641 self.mutation_count += 1;
1642 }
1643 Mutation::InsertEdge {
1644 src_vid,
1645 dst_vid,
1646 edge_type,
1647 eid,
1648 version: _,
1649 properties,
1650 edge_type_name,
1651 } => {
1652 self.current_version += 1;
1653 match self.apply_edge_insertion(src_vid, dst_vid, edge_type, eid, properties) {
1658 Ok(()) => {
1659 if let Some(name) = edge_type_name {
1661 self.edge_types.insert(eid, name);
1662 }
1663 }
1664 Err(e) => {
1665 tracing::warn!(
1666 ?eid,
1667 ?src_vid,
1668 ?dst_vid,
1669 error = %e,
1670 "WAL replay: skipping edge insertion to a deleted endpoint (issue #77)"
1671 );
1672 }
1673 }
1674 }
1675 Mutation::DeleteEdge {
1676 eid,
1677 src_vid,
1678 dst_vid,
1679 edge_type,
1680 version: _,
1681 } => {
1682 self.current_version += 1;
1683 self.apply_edge_deletion(eid, src_vid, dst_vid, edge_type);
1684 }
1685 }
1686 }
1687 Ok(())
1688 }
1689}
1690
1691#[cfg(test)]
1692mod tests {
1693 use super::*;
1694
1695 #[test]
1696 fn test_l0_buffer_ops() -> Result<()> {
1697 let mut l0 = L0Buffer::new(0, None);
1698 let vid_a = Vid::new(1);
1699 let vid_b = Vid::new(2);
1700 let eid_ab = Eid::new(101);
1701
1702 l0.insert_edge(vid_a, vid_b, 1, eid_ab, HashMap::new(), None)?;
1703
1704 let neighbors = l0.get_neighbors(vid_a, 1, Direction::Outgoing);
1705 assert_eq!(neighbors.len(), 1);
1706 assert_eq!(neighbors[0].0, vid_b);
1707 assert_eq!(neighbors[0].1, eid_ab);
1708
1709 l0.delete_edge(eid_ab, vid_a, vid_b, 1)?;
1710 assert!(l0.is_tombstoned(eid_ab));
1711
1712 let neighbors_after = l0.get_neighbors(vid_a, 1, Direction::Outgoing);
1714 assert_eq!(neighbors_after.len(), 0);
1715
1716 Ok(())
1717 }
1718
1719 #[test]
1725 fn validate_merge_rejects_edge_to_tombstoned_endpoint() {
1726 let mut main = L0Buffer::new(0, None);
1727 let vid_a = Vid::new(1);
1728 let vid_b = Vid::new(2);
1729 main.insert_vertex(vid_a, HashMap::new());
1730 main.insert_vertex(vid_b, HashMap::new());
1731 main.delete_vertex(vid_b).unwrap(); let mut tx = L0Buffer::new(0, None);
1735 let eid = Eid::new(101);
1736 tx.insert_edge(vid_a, vid_b, 1, eid, HashMap::new(), None)
1737 .unwrap();
1738
1739 assert!(
1740 main.validate_merge_edge_endpoints(&tx).is_err(),
1741 "edge to a tombstoned endpoint must be rejected before merge"
1742 );
1743 assert!(
1745 main.merge(&tx).is_err(),
1746 "merge must reject, not bail mid-apply"
1747 );
1748 assert!(
1749 !main.edge_endpoints.contains_key(&eid),
1750 "a rejected merge must not have partially applied the edge"
1751 );
1752 }
1753
1754 #[test]
1757 fn validate_merge_allows_edge_when_endpoint_reinserted() {
1758 let mut main = L0Buffer::new(0, None);
1759 let vid_a = Vid::new(1);
1760 let vid_b = Vid::new(2);
1761 main.insert_vertex(vid_a, HashMap::new());
1762 main.insert_vertex(vid_b, HashMap::new());
1763 main.delete_vertex(vid_b).unwrap();
1764
1765 let mut tx = L0Buffer::new(0, None);
1766 tx.insert_vertex(vid_b, HashMap::new()); let eid = Eid::new(101);
1768 tx.insert_edge(vid_a, vid_b, 1, eid, HashMap::new(), None)
1769 .unwrap();
1770
1771 assert!(main.validate_merge_edge_endpoints(&tx).is_ok());
1772 assert!(main.merge(&tx).is_ok());
1773 assert!(main.edge_endpoints.contains_key(&eid));
1774 }
1775
1776 #[test]
1778 fn validate_merge_allows_edge_to_live_endpoints() {
1779 let mut main = L0Buffer::new(0, None);
1780 let vid_a = Vid::new(1);
1781 let vid_b = Vid::new(2);
1782 main.insert_vertex(vid_a, HashMap::new());
1783 main.insert_vertex(vid_b, HashMap::new());
1784
1785 let mut tx = L0Buffer::new(0, None);
1786 let eid = Eid::new(101);
1787 tx.insert_edge(vid_a, vid_b, 1, eid, HashMap::new(), None)
1788 .unwrap();
1789
1790 assert!(main.validate_merge_edge_endpoints(&tx).is_ok());
1791 assert!(main.merge(&tx).is_ok());
1792 assert!(main.edge_endpoints.contains_key(&eid));
1793 }
1794
1795 #[test]
1796 fn test_l0_buffer_multiple_edges() -> Result<()> {
1797 let mut l0 = L0Buffer::new(0, None);
1798 let vid_a = Vid::new(1);
1799 let vid_b = Vid::new(2);
1800 let vid_c = Vid::new(3);
1801 let eid_ab = Eid::new(101);
1802 let eid_ac = Eid::new(102);
1803
1804 l0.insert_edge(vid_a, vid_b, 1, eid_ab, HashMap::new(), None)?;
1805 l0.insert_edge(vid_a, vid_c, 1, eid_ac, HashMap::new(), None)?;
1806
1807 let neighbors = l0.get_neighbors(vid_a, 1, Direction::Outgoing);
1808 assert_eq!(neighbors.len(), 2);
1809
1810 l0.delete_edge(eid_ab, vid_a, vid_b, 1)?;
1812
1813 let neighbors_after = l0.get_neighbors(vid_a, 1, Direction::Outgoing);
1815 assert_eq!(neighbors_after.len(), 1);
1816 assert_eq!(neighbors_after[0].0, vid_c);
1817
1818 Ok(())
1819 }
1820
1821 #[test]
1822 fn test_l0_buffer_edge_type_filter() -> Result<()> {
1823 let mut l0 = L0Buffer::new(0, None);
1824 let vid_a = Vid::new(1);
1825 let vid_b = Vid::new(2);
1826 let vid_c = Vid::new(3);
1827 let eid_ab = Eid::new(101);
1828 let eid_ac = Eid::new(201); l0.insert_edge(vid_a, vid_b, 1, eid_ab, HashMap::new(), None)?;
1831 l0.insert_edge(vid_a, vid_c, 2, eid_ac, HashMap::new(), None)?;
1832
1833 let type1_neighbors = l0.get_neighbors(vid_a, 1, Direction::Outgoing);
1835 assert_eq!(type1_neighbors.len(), 1);
1836 assert_eq!(type1_neighbors[0].0, vid_b);
1837
1838 let type2_neighbors = l0.get_neighbors(vid_a, 2, Direction::Outgoing);
1840 assert_eq!(type2_neighbors.len(), 1);
1841 assert_eq!(type2_neighbors[0].0, vid_c);
1842
1843 Ok(())
1844 }
1845
1846 #[test]
1847 fn test_l0_buffer_incoming_edges() -> Result<()> {
1848 let mut l0 = L0Buffer::new(0, None);
1849 let vid_a = Vid::new(1);
1850 let vid_b = Vid::new(2);
1851 let vid_c = Vid::new(3);
1852 let eid_ab = Eid::new(101);
1853 let eid_cb = Eid::new(102);
1854
1855 l0.insert_edge(vid_a, vid_b, 1, eid_ab, HashMap::new(), None)?;
1857 l0.insert_edge(vid_c, vid_b, 1, eid_cb, HashMap::new(), None)?;
1858
1859 let incoming = l0.get_neighbors(vid_b, 1, Direction::Incoming);
1861 assert_eq!(incoming.len(), 2);
1862
1863 Ok(())
1864 }
1865
1866 #[test]
1868 fn test_merge_empty_props_edge() -> Result<()> {
1869 let mut main_l0 = L0Buffer::new(0, None);
1870 let mut tx_l0 = L0Buffer::new(0, None);
1871
1872 let vid_a = Vid::new(1);
1873 let vid_b = Vid::new(2);
1874 let eid_ab = Eid::new(101);
1875
1876 tx_l0.insert_edge(vid_a, vid_b, 1, eid_ab, HashMap::new(), None)?;
1878
1879 assert!(tx_l0.edge_endpoints.contains_key(&eid_ab));
1881 assert!(!tx_l0.edge_properties.contains_key(&eid_ab)); main_l0.merge(&tx_l0)?;
1885
1886 assert!(main_l0.edge_endpoints.contains_key(&eid_ab));
1888 let neighbors = main_l0.get_neighbors(vid_a, 1, Direction::Outgoing);
1889 assert_eq!(neighbors.len(), 1);
1890 assert_eq!(neighbors[0].0, vid_b);
1891
1892 Ok(())
1893 }
1894
1895 #[test]
1897 fn test_replay_crdt_merge() -> Result<()> {
1898 use crate::runtime::wal::Mutation;
1899 use serde_json::json;
1900 use uni_common::Value;
1901
1902 let mut l0 = L0Buffer::new(0, None);
1903 let vid = Vid::new(1);
1904
1905 let counter1: Value = json!({
1908 "t": "gc",
1909 "d": {"counts": {"node1": 5}}
1910 })
1911 .into();
1912 let counter2: Value = json!({
1913 "t": "gc",
1914 "d": {"counts": {"node2": 3}}
1915 })
1916 .into();
1917
1918 let mut props1 = HashMap::new();
1920 props1.insert("counter".to_string(), counter1.clone());
1921 l0.replay_mutations(vec![Mutation::InsertVertex {
1922 vid,
1923 properties: props1,
1924 labels: vec![],
1925 }])?;
1926
1927 let mut props2 = HashMap::new();
1929 props2.insert("counter".to_string(), counter2.clone());
1930 l0.replay_mutations(vec![Mutation::InsertVertex {
1931 vid,
1932 properties: props2,
1933 labels: vec![],
1934 }])?;
1935
1936 let stored_props = l0.vertex_properties.get(&vid).unwrap();
1938 let stored_counter = stored_props.get("counter").unwrap();
1939
1940 let stored_json: serde_json::Value = stored_counter.clone().into();
1942 let data = stored_json.get("d").unwrap();
1944 let counts = data.get("counts").unwrap();
1945 assert_eq!(counts.get("node1"), Some(&json!(5)));
1946 assert_eq!(counts.get("node2"), Some(&json!(3)));
1947
1948 Ok(())
1949 }
1950
1951 #[test]
1952 fn test_merge_preserves_vertex_timestamps() -> Result<()> {
1953 let mut l0_main = L0Buffer::new(0, None);
1954 let mut l0_tx = L0Buffer::new(0, None);
1955 let vid = Vid::new(1);
1956
1957 let ts_main_created = 1000;
1959 let ts_main_updated = 1100;
1960 l0_main.insert_vertex(vid, HashMap::new());
1961 l0_main.vertex_created_at.insert(vid, ts_main_created);
1962 l0_main.vertex_updated_at.insert(vid, ts_main_updated);
1963
1964 let ts_tx_created = 2000; let ts_tx_updated = 2100; l0_tx.insert_vertex(vid, HashMap::new());
1968 l0_tx.vertex_created_at.insert(vid, ts_tx_created);
1969 l0_tx.vertex_updated_at.insert(vid, ts_tx_updated);
1970
1971 l0_main.merge(&l0_tx)?;
1973
1974 assert_eq!(
1976 *l0_main.vertex_created_at.get(&vid).unwrap(),
1977 ts_main_created,
1978 "created_at should preserve oldest timestamp"
1979 );
1980
1981 assert_eq!(
1983 *l0_main.vertex_updated_at.get(&vid).unwrap(),
1984 ts_tx_updated,
1985 "updated_at should use latest timestamp"
1986 );
1987
1988 Ok(())
1989 }
1990
1991 #[test]
1992 fn test_merge_preserves_edge_timestamps() -> Result<()> {
1993 let mut l0_main = L0Buffer::new(0, None);
1994 let mut l0_tx = L0Buffer::new(0, None);
1995 let vid_a = Vid::new(1);
1996 let vid_b = Vid::new(2);
1997 let eid = Eid::new(100);
1998
1999 let ts_main_created = 1000;
2001 let ts_main_updated = 1100;
2002 l0_main.insert_edge(vid_a, vid_b, 1, eid, HashMap::new(), None)?;
2003 l0_main.edge_created_at.insert(eid, ts_main_created);
2004 l0_main.edge_updated_at.insert(eid, ts_main_updated);
2005
2006 let ts_tx_created = 2000; let ts_tx_updated = 2100; l0_tx.insert_edge(vid_a, vid_b, 1, eid, HashMap::new(), None)?;
2010 l0_tx.edge_created_at.insert(eid, ts_tx_created);
2011 l0_tx.edge_updated_at.insert(eid, ts_tx_updated);
2012
2013 l0_main.merge(&l0_tx)?;
2015
2016 assert_eq!(
2018 *l0_main.edge_created_at.get(&eid).unwrap(),
2019 ts_main_created,
2020 "edge created_at should preserve oldest timestamp"
2021 );
2022
2023 assert_eq!(
2025 *l0_main.edge_updated_at.get(&eid).unwrap(),
2026 ts_tx_updated,
2027 "edge updated_at should use latest timestamp"
2028 );
2029
2030 Ok(())
2031 }
2032
2033 #[test]
2034 fn test_merge_created_at_not_overwritten_for_existing_vertex() -> Result<()> {
2035 use uni_common::Value;
2036
2037 let mut l0_main = L0Buffer::new(0, None);
2038 let mut l0_tx = L0Buffer::new(0, None);
2039 let vid = Vid::new(1);
2040
2041 let ts_original = 1000;
2043 l0_main.insert_vertex(vid, HashMap::new());
2044 l0_main.vertex_created_at.insert(vid, ts_original);
2045 l0_main.vertex_updated_at.insert(vid, ts_original);
2046
2047 let ts_tx = 2000;
2049 let mut props = HashMap::new();
2050 props.insert("updated".to_string(), Value::String("yes".to_string()));
2051 l0_tx.insert_vertex(vid, props);
2052 l0_tx.vertex_created_at.insert(vid, ts_tx);
2053 l0_tx.vertex_updated_at.insert(vid, ts_tx);
2054
2055 l0_main.merge(&l0_tx)?;
2057
2058 assert_eq!(
2060 *l0_main.vertex_created_at.get(&vid).unwrap(),
2061 ts_original,
2062 "created_at must not be overwritten for existing vertex"
2063 );
2064
2065 assert_eq!(
2067 *l0_main.vertex_updated_at.get(&vid).unwrap(),
2068 ts_tx,
2069 "updated_at should reflect transaction timestamp"
2070 );
2071
2072 assert!(
2074 l0_main
2075 .vertex_properties
2076 .get(&vid)
2077 .unwrap()
2078 .contains_key("updated")
2079 );
2080
2081 Ok(())
2082 }
2083
2084 #[test]
2086 fn test_replay_mutations_preserves_vertex_labels() -> Result<()> {
2087 use crate::runtime::wal::Mutation;
2088
2089 let mut l0 = L0Buffer::new(0, None);
2090 let vid = Vid::new(42);
2091
2092 let mutations = vec![Mutation::InsertVertex {
2094 vid,
2095 properties: {
2096 let mut props = HashMap::new();
2097 props.insert(
2098 "name".to_string(),
2099 uni_common::Value::String("Alice".to_string()),
2100 );
2101 props
2102 },
2103 labels: vec!["Person".to_string(), "User".to_string()],
2104 }];
2105
2106 l0.replay_mutations(mutations)?;
2108
2109 assert!(l0.vertex_properties.contains_key(&vid));
2111
2112 let labels = l0.get_vertex_labels(vid).expect("Labels should exist");
2114 assert_eq!(labels.len(), 2);
2115 assert!(labels.contains(&"Person".to_string()));
2116 assert!(labels.contains(&"User".to_string()));
2117
2118 let person_vids = l0.vids_for_label("Person");
2120 assert_eq!(person_vids.len(), 1);
2121 assert_eq!(person_vids[0], vid);
2122
2123 let user_vids = l0.vids_for_label("User");
2124 assert_eq!(user_vids.len(), 1);
2125 assert_eq!(user_vids[0], vid);
2126
2127 Ok(())
2128 }
2129
2130 #[test]
2132 fn test_replay_mutations_preserves_delete_vertex_labels() -> Result<()> {
2133 use crate::runtime::wal::Mutation;
2134
2135 let mut l0 = L0Buffer::new(0, None);
2136 let vid = Vid::new(99);
2137
2138 l0.insert_vertex_with_labels(
2140 vid,
2141 HashMap::new(),
2142 &["Person".to_string(), "Admin".to_string()],
2143 );
2144
2145 assert!(l0.vertex_properties.contains_key(&vid));
2147 let labels = l0.get_vertex_labels(vid).expect("Labels should exist");
2148 assert_eq!(labels.len(), 2);
2149
2150 let mutations = vec![Mutation::DeleteVertex {
2152 vid,
2153 labels: vec!["Person".to_string(), "Admin".to_string()],
2154 }];
2155
2156 l0.replay_mutations(mutations)?;
2158
2159 assert!(l0.vertex_tombstones.contains(&vid));
2161
2162 let labels = l0.get_vertex_labels(vid);
2165 assert!(
2166 labels.is_some(),
2167 "Labels should be preserved even after deletion for tombstone flushing"
2168 );
2169
2170 Ok(())
2171 }
2172
2173 #[test]
2175 fn test_replay_mutations_preserves_edge_type_name() -> Result<()> {
2176 use crate::runtime::wal::Mutation;
2177
2178 let mut l0 = L0Buffer::new(0, None);
2179 let src = Vid::new(1);
2180 let dst = Vid::new(2);
2181 let eid = Eid::new(500);
2182 let edge_type = 100;
2183
2184 let mutations = vec![Mutation::InsertEdge {
2186 src_vid: src,
2187 dst_vid: dst,
2188 edge_type,
2189 eid,
2190 version: 1,
2191 properties: {
2192 let mut props = HashMap::new();
2193 props.insert("since".to_string(), uni_common::Value::Int(2020));
2194 props
2195 },
2196 edge_type_name: Some("KNOWS".to_string()),
2197 }];
2198
2199 l0.replay_mutations(mutations)?;
2201
2202 assert!(l0.edge_endpoints.contains_key(&eid));
2204
2205 let type_name = l0.get_edge_type(eid).expect("Edge type name should exist");
2207 assert_eq!(type_name, "KNOWS");
2208
2209 let knows_eids = l0.eids_for_type("KNOWS");
2211 assert_eq!(knows_eids.len(), 1);
2212 assert_eq!(knows_eids[0], eid);
2213
2214 Ok(())
2215 }
2216
2217 #[test]
2219 fn test_edge_type_mapping_survives_multiple_replays() -> Result<()> {
2220 use crate::runtime::wal::Mutation;
2221
2222 let mut l0 = L0Buffer::new(0, None);
2223
2224 let mutations = vec![
2226 Mutation::InsertEdge {
2227 src_vid: Vid::new(1),
2228 dst_vid: Vid::new(2),
2229 edge_type: 100,
2230 eid: Eid::new(1000),
2231 version: 1,
2232 properties: HashMap::new(),
2233 edge_type_name: Some("KNOWS".to_string()),
2234 },
2235 Mutation::InsertEdge {
2236 src_vid: Vid::new(2),
2237 dst_vid: Vid::new(3),
2238 edge_type: 101,
2239 eid: Eid::new(1001),
2240 version: 2,
2241 properties: HashMap::new(),
2242 edge_type_name: Some("LIKES".to_string()),
2243 },
2244 Mutation::InsertEdge {
2245 src_vid: Vid::new(3),
2246 dst_vid: Vid::new(1),
2247 edge_type: 100,
2248 eid: Eid::new(1002),
2249 version: 3,
2250 properties: HashMap::new(),
2251 edge_type_name: Some("KNOWS".to_string()),
2252 },
2253 ];
2254
2255 l0.replay_mutations(mutations)?;
2256
2257 assert_eq!(l0.get_edge_type(Eid::new(1000)), Some("KNOWS"));
2259 assert_eq!(l0.get_edge_type(Eid::new(1001)), Some("LIKES"));
2260 assert_eq!(l0.get_edge_type(Eid::new(1002)), Some("KNOWS"));
2261
2262 let knows_edges = l0.eids_for_type("KNOWS");
2264 assert_eq!(knows_edges.len(), 2);
2265 assert!(knows_edges.contains(&Eid::new(1000)));
2266 assert!(knows_edges.contains(&Eid::new(1002)));
2267
2268 let likes_edges = l0.eids_for_type("LIKES");
2269 assert_eq!(likes_edges.len(), 1);
2270 assert_eq!(likes_edges[0], Eid::new(1001));
2271
2272 Ok(())
2273 }
2274
2275 #[test]
2277 fn test_replay_mutations_combined_labels_and_edge_types() -> Result<()> {
2278 use crate::runtime::wal::Mutation;
2279
2280 let mut l0 = L0Buffer::new(0, None);
2281 let alice = Vid::new(1);
2282 let bob = Vid::new(2);
2283 let eid = Eid::new(100);
2284
2285 let mutations = vec![
2287 Mutation::InsertVertex {
2289 vid: alice,
2290 properties: {
2291 let mut props = HashMap::new();
2292 props.insert(
2293 "name".to_string(),
2294 uni_common::Value::String("Alice".to_string()),
2295 );
2296 props
2297 },
2298 labels: vec!["Person".to_string()],
2299 },
2300 Mutation::InsertVertex {
2302 vid: bob,
2303 properties: {
2304 let mut props = HashMap::new();
2305 props.insert(
2306 "name".to_string(),
2307 uni_common::Value::String("Bob".to_string()),
2308 );
2309 props
2310 },
2311 labels: vec!["Person".to_string()],
2312 },
2313 Mutation::InsertEdge {
2315 src_vid: alice,
2316 dst_vid: bob,
2317 edge_type: 1,
2318 eid,
2319 version: 3,
2320 properties: HashMap::new(),
2321 edge_type_name: Some("KNOWS".to_string()),
2322 },
2323 ];
2324
2325 l0.replay_mutations(mutations)?;
2327
2328 assert_eq!(l0.get_vertex_labels(alice).unwrap().len(), 1);
2330 assert_eq!(l0.get_vertex_labels(bob).unwrap().len(), 1);
2331 assert_eq!(l0.vids_for_label("Person").len(), 2);
2332
2333 assert_eq!(l0.get_edge_type(eid).unwrap(), "KNOWS");
2335 assert_eq!(l0.eids_for_type("KNOWS").len(), 1);
2336
2337 let alice_neighbors = l0.get_neighbors(alice, 1, Direction::Outgoing);
2339 assert_eq!(alice_neighbors.len(), 1);
2340 assert_eq!(alice_neighbors[0].0, bob);
2341
2342 Ok(())
2343 }
2344
2345 #[test]
2347 fn test_replay_mutations_backward_compat_empty_labels() -> Result<()> {
2348 use crate::runtime::wal::Mutation;
2349
2350 let mut l0 = L0Buffer::new(0, None);
2351 let vid = Vid::new(1);
2352
2353 let mutations = vec![Mutation::InsertVertex {
2356 vid,
2357 properties: HashMap::new(),
2358 labels: vec![], }];
2360
2361 l0.replay_mutations(mutations)?;
2362
2363 assert!(l0.vertex_properties.contains_key(&vid));
2365
2366 let labels = l0.get_vertex_labels(vid);
2368 assert!(labels.is_some(), "Labels entry should exist even if empty");
2369 assert_eq!(labels.unwrap().len(), 0);
2370
2371 Ok(())
2372 }
2373
2374 #[test]
2375 fn test_now_nanos_returns_nanosecond_range() {
2376 let now = now_nanos();
2380
2381 assert!(
2383 now > 1_700_000_000_000_000_000,
2384 "now_nanos() returned {}, expected > 1.7e18 for nanoseconds",
2385 now
2386 );
2387
2388 assert!(
2390 now < 4_100_000_000_000_000_000,
2391 "now_nanos() returned {}, expected < 4.1e18",
2392 now
2393 );
2394 }
2395}