1#![allow(dead_code)]
7
8use crate::model::{Triple, TriplePattern};
9use crate::OxirsError;
10use scirs2_core::random::{Random, RngExt};
11use serde::{Deserialize, Serialize};
12use std::collections::{BTreeMap, BTreeSet, HashMap};
13use std::sync::Arc;
14use tokio::sync::RwLock;
15
16#[derive(Debug, Clone)]
18pub struct CrdtConfig {
19 pub node_id: String,
21 pub crdt_type: CrdtType,
23 pub gc_config: GcConfig,
25 pub delta_config: DeltaConfig,
27}
28
29#[derive(Debug, Clone)]
31pub enum CrdtType {
32 GSet,
34 TwoPhaseSet,
36 AddRemovePartialOrder,
38 OrSet,
40 LwwSet,
42 MvRegister,
44 RdfCrdt,
46}
47
48#[derive(Debug, Clone)]
50pub struct GcConfig {
51 pub auto_gc: bool,
53 pub interval_secs: u64,
55 pub tombstone_ttl_secs: u64,
57 pub batch_size: usize,
59}
60
61impl Default for GcConfig {
62 fn default() -> Self {
63 GcConfig {
64 auto_gc: true,
65 interval_secs: 3600, tombstone_ttl_secs: 86400 * 7, batch_size: 1000,
68 }
69 }
70}
71
72#[derive(Debug, Clone)]
74pub struct DeltaConfig {
75 pub enabled: bool,
77 pub max_delta_size: usize,
79 pub buffer_size: usize,
81 pub compression: bool,
83}
84
85impl Default for DeltaConfig {
86 fn default() -> Self {
87 DeltaConfig {
88 enabled: true,
89 max_delta_size: 10000,
90 buffer_size: 100000,
91 compression: true,
92 }
93 }
94}
95
96pub trait Crdt: Send + Sync {
98 type Delta: Send + Sync + Clone + Serialize + for<'de> Deserialize<'de>;
100
101 fn merge(&mut self, other: &Self);
103
104 fn delta(&self) -> Option<Self::Delta>;
106
107 fn apply_delta(&mut self, delta: Self::Delta);
109
110 fn reset_delta(&mut self);
112}
113
114#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
116pub struct ElementId {
117 pub timestamp: u64,
119 pub node_id: String,
121 pub random: u64,
123}
124
125impl ElementId {
126 pub fn new(timestamp: u64, node_id: String) -> Self {
128 ElementId {
129 timestamp,
130 node_id,
131 random: {
132 let mut rng = Random::default();
133 rng.random::<u64>()
134 },
135 }
136 }
137}
138
139#[derive(Debug, Clone)]
141pub struct GrowSet<T: Clone + Ord + Send + Sync> {
142 elements: BTreeSet<T>,
144 delta_elements: Option<BTreeSet<T>>,
146}
147
148impl<T: Clone + Ord + Send + Sync + Serialize + for<'de> Deserialize<'de>> Default for GrowSet<T> {
149 fn default() -> Self {
150 Self::new()
151 }
152}
153
154impl<T: Clone + Ord + Send + Sync + Serialize + for<'de> Deserialize<'de>> GrowSet<T> {
155 pub fn new() -> Self {
157 GrowSet {
158 elements: BTreeSet::new(),
159 delta_elements: Some(BTreeSet::new()),
160 }
161 }
162
163 pub fn add(&mut self, element: T) {
165 if self.elements.insert(element.clone()) {
166 if let Some(ref mut delta) = self.delta_elements {
167 delta.insert(element);
168 }
169 }
170 }
171
172 pub fn contains(&self, element: &T) -> bool {
174 self.elements.contains(element)
175 }
176
177 pub fn elements(&self) -> &BTreeSet<T> {
179 &self.elements
180 }
181}
182
183impl<T: Clone + Ord + Send + Sync + Serialize + for<'de> Deserialize<'de>> Crdt for GrowSet<T> {
184 type Delta = BTreeSet<T>;
185
186 fn merge(&mut self, other: &Self) {
187 for element in &other.elements {
188 self.add(element.clone());
189 }
190 }
191
192 fn delta(&self) -> Option<Self::Delta> {
193 self.delta_elements.clone()
194 }
195
196 fn apply_delta(&mut self, delta: Self::Delta) {
197 for element in delta {
198 self.elements.insert(element);
199 }
200 }
201
202 fn reset_delta(&mut self) {
203 self.delta_elements = Some(BTreeSet::new());
204 }
205}
206
207#[derive(Debug, Clone)]
209pub struct TwoPhaseSet<T: Clone + Ord + Send + Sync> {
210 added: BTreeSet<T>,
212 removed: BTreeSet<T>,
214 delta_added: Option<BTreeSet<T>>,
216 delta_removed: Option<BTreeSet<T>>,
217}
218
219impl<T: Clone + Ord + Send + Sync + Serialize + for<'de> Deserialize<'de>> Default
220 for TwoPhaseSet<T>
221{
222 fn default() -> Self {
223 Self::new()
224 }
225}
226
227impl<T: Clone + Ord + Send + Sync + Serialize + for<'de> Deserialize<'de>> TwoPhaseSet<T> {
228 pub fn new() -> Self {
230 TwoPhaseSet {
231 added: BTreeSet::new(),
232 removed: BTreeSet::new(),
233 delta_added: Some(BTreeSet::new()),
234 delta_removed: Some(BTreeSet::new()),
235 }
236 }
237
238 pub fn add(&mut self, element: T) {
240 if !self.removed.contains(&element) && self.added.insert(element.clone()) {
241 if let Some(ref mut delta) = self.delta_added {
242 delta.insert(element);
243 }
244 }
245 }
246
247 pub fn remove(&mut self, element: T) {
249 if self.added.contains(&element) && self.removed.insert(element.clone()) {
250 if let Some(ref mut delta) = self.delta_removed {
251 delta.insert(element);
252 }
253 }
254 }
255
256 pub fn contains(&self, element: &T) -> bool {
258 self.added.contains(element) && !self.removed.contains(element)
259 }
260
261 pub fn elements(&self) -> BTreeSet<T> {
263 self.added.difference(&self.removed).cloned().collect()
264 }
265}
266
267#[derive(Debug, Clone)]
269pub struct OrSet<T: Clone + Ord + Send + Sync> {
270 elements: BTreeMap<T, BTreeSet<ElementId>>,
272 tombstones: BTreeMap<T, BTreeSet<ElementId>>,
274 node_id: String,
276 clock: u64,
278 delta: Option<OrSetDelta<T>>,
280}
281
282#[derive(Debug, Clone, Serialize, Deserialize)]
283pub struct OrSetDelta<T: Clone + Ord> {
284 added: BTreeMap<T, BTreeSet<ElementId>>,
286 removed: BTreeMap<T, BTreeSet<ElementId>>,
288}
289
290impl<T: Clone + Ord + Send + Sync + Serialize + for<'de> Deserialize<'de>> OrSet<T> {
291 pub fn new(node_id: String) -> Self {
293 OrSet {
294 elements: BTreeMap::new(),
295 tombstones: BTreeMap::new(),
296 node_id,
297 clock: 0,
298 delta: Some(OrSetDelta {
299 added: BTreeMap::new(),
300 removed: BTreeMap::new(),
301 }),
302 }
303 }
304
305 pub fn add(&mut self, element: T) {
307 self.clock += 1;
308 let tag = ElementId::new(self.clock, self.node_id.clone());
309
310 self.elements
311 .entry(element.clone())
312 .or_default()
313 .insert(tag.clone());
314
315 if let Some(ref mut delta) = self.delta {
316 delta
317 .added
318 .entry(element)
319 .or_insert_with(BTreeSet::new)
320 .insert(tag);
321 }
322 }
323
324 pub fn remove(&mut self, element: &T) {
326 if let Some(tags) = self.elements.get(element).cloned() {
327 self.tombstones.insert(element.clone(), tags.clone());
328 self.elements.remove(element);
329
330 if let Some(ref mut delta) = self.delta {
331 delta.removed.insert(element.clone(), tags);
332 }
333 }
334 }
335
336 pub fn contains(&self, element: &T) -> bool {
338 if let Some(tags) = self.elements.get(element) {
339 if let Some(tombstone_tags) = self.tombstones.get(element) {
340 !tags.is_subset(tombstone_tags)
342 } else {
343 true
344 }
345 } else {
346 false
347 }
348 }
349
350 pub fn elements(&self) -> BTreeSet<T> {
352 self.elements
353 .keys()
354 .filter(|e| self.contains(e))
355 .cloned()
356 .collect()
357 }
358}
359
360impl<T: Clone + Ord + Send + Sync + Serialize + for<'de> Deserialize<'de>> Crdt for OrSet<T> {
361 type Delta = OrSetDelta<T>;
362
363 fn merge(&mut self, other: &Self) {
364 for (element, tags) in &other.elements {
366 self.elements
367 .entry(element.clone())
368 .or_default()
369 .extend(tags.iter().cloned());
370 }
371
372 for (element, tags) in &other.tombstones {
374 self.tombstones
375 .entry(element.clone())
376 .or_default()
377 .extend(tags.iter().cloned());
378 }
379
380 let to_remove: Vec<_> = self
382 .elements
383 .iter()
384 .filter(|(e, tags)| {
385 if let Some(tombstone_tags) = self.tombstones.get(e) {
386 tags.is_subset(tombstone_tags)
387 } else {
388 false
389 }
390 })
391 .map(|(e, _)| e.clone())
392 .collect();
393
394 for element in to_remove {
395 self.elements.remove(&element);
396 }
397
398 self.clock = self.clock.max(other.clock);
400 }
401
402 fn delta(&self) -> Option<Self::Delta> {
403 self.delta.clone()
404 }
405
406 fn apply_delta(&mut self, delta: Self::Delta) {
407 for (element, tags) in delta.added {
409 self.elements.entry(element).or_default().extend(tags);
410 }
411
412 for (element, tags) in delta.removed {
414 self.tombstones
415 .entry(element.clone())
416 .or_default()
417 .extend(tags);
418
419 if let Some(elem_tags) = self.elements.get(&element) {
421 if let Some(tombstone_tags) = self.tombstones.get(&element) {
422 if elem_tags.is_subset(tombstone_tags) {
423 self.elements.remove(&element);
424 }
425 }
426 }
427 }
428 }
429
430 fn reset_delta(&mut self) {
431 self.delta = Some(OrSetDelta {
432 added: BTreeMap::new(),
433 removed: BTreeMap::new(),
434 });
435 }
436}
437
438pub struct RdfCrdt {
440 config: CrdtConfig,
442 triples: OrSet<Triple>,
444 predicate_index: HashMap<String, OrSet<Triple>>,
446 subject_index: HashMap<String, OrSet<Triple>>,
448 stats: Arc<RwLock<CrdtStats>>,
450}
451
452#[derive(Debug, Default)]
454struct CrdtStats {
455 total_ops: u64,
457 add_ops: u64,
459 remove_ops: u64,
461 merge_ops: u64,
463 triple_count: usize,
465 #[allow(dead_code)]
467 tombstone_count: usize,
468}
469
470fn subject_index_key(subject: &crate::model::Subject) -> String {
482 match subject {
483 crate::model::Subject::NamedNode(nn) => nn.as_str().to_string(),
484 crate::model::Subject::BlankNode(bn) => bn.as_str().to_string(),
485 crate::model::Subject::Variable(v) => v.as_str().to_string(),
486 crate::model::Subject::QuotedTriple(qt) => format!("{qt}"),
487 }
488}
489
490impl RdfCrdt {
491 pub async fn new(config: CrdtConfig) -> Result<Self, OxirsError> {
493 let node_id = config.node_id.clone();
494
495 Ok(RdfCrdt {
496 config,
497 triples: OrSet::new(node_id),
498 predicate_index: HashMap::new(),
499 subject_index: HashMap::new(),
500 stats: Arc::new(RwLock::new(CrdtStats::default())),
501 })
502 }
503
504 pub async fn add_triple(&mut self, triple: Triple) -> Result<(), OxirsError> {
506 self.triples.add(triple.clone());
508
509 let predicate_str = match triple.predicate() {
511 crate::model::Predicate::NamedNode(nn) => nn.as_str(),
512 crate::model::Predicate::Variable(v) => v.as_str(),
513 };
514 self.predicate_index
515 .entry(predicate_str.to_string())
516 .or_insert_with(|| OrSet::new(self.config.node_id.clone()))
517 .add(triple.clone());
518
519 let subject_key = subject_index_key(triple.subject());
521 self.subject_index
522 .entry(subject_key)
523 .or_insert_with(|| OrSet::new(self.config.node_id.clone()))
524 .add(triple);
525
526 let mut stats = self.stats.write().await;
528 stats.total_ops += 1;
529 stats.add_ops += 1;
530 stats.triple_count = self.triples.elements().len();
531
532 Ok(())
533 }
534
535 pub async fn remove_triple(&mut self, triple: &Triple) -> Result<(), OxirsError> {
537 self.triples.remove(triple);
539
540 let predicate_str = match triple.predicate() {
542 crate::model::Predicate::NamedNode(nn) => nn.as_str(),
543 crate::model::Predicate::Variable(v) => v.as_str(),
544 };
545 if let Some(predicate_set) = self.predicate_index.get_mut(predicate_str) {
546 predicate_set.remove(triple);
547 }
548
549 let subject_key = subject_index_key(triple.subject());
551 if let Some(subject_set) = self.subject_index.get_mut(&subject_key) {
552 subject_set.remove(triple);
553 }
554
555 let mut stats = self.stats.write().await;
557 stats.total_ops += 1;
558 stats.remove_ops += 1;
559 stats.triple_count = self.triples.elements().len();
560
561 Ok(())
562 }
563
564 pub async fn query(&self, pattern: &TriplePattern) -> Result<Vec<Triple>, OxirsError> {
566 let is_quoted_triple_subject = matches!(
576 pattern.subject(),
577 Some(crate::model::SubjectPattern::QuotedTriple(_))
578 );
579
580 let results = match (pattern.subject(), pattern.predicate(), pattern.object()) {
581 (Some(_), Some(_), _) | (Some(_), None, _) if is_quoted_triple_subject => self
582 .triples
583 .elements()
584 .into_iter()
585 .filter(|t| pattern.matches(t))
586 .collect(),
587 (Some(subject), Some(_predicate), _) => {
588 if let Some(subject_set) = self.subject_index.get(subject.as_str()) {
590 subject_set
591 .elements()
592 .into_iter()
593 .filter(|t| pattern.matches(t))
594 .collect()
595 } else {
596 Vec::new()
597 }
598 }
599 (Some(subject), None, _) => {
600 if let Some(subject_set) = self.subject_index.get(subject.as_str()) {
602 subject_set
603 .elements()
604 .into_iter()
605 .filter(|t| pattern.matches(t))
606 .collect()
607 } else {
608 Vec::new()
609 }
610 }
611 (None, Some(predicate), _) => {
612 if let Some(predicate_set) = self.predicate_index.get(predicate.as_str()) {
614 predicate_set
615 .elements()
616 .into_iter()
617 .filter(|t| pattern.matches(t))
618 .collect()
619 } else {
620 Vec::new()
621 }
622 }
623 _ => {
624 self.triples
626 .elements()
627 .into_iter()
628 .filter(|t| pattern.matches(t))
629 .collect()
630 }
631 };
632
633 Ok(results)
634 }
635
636 pub async fn merge(&mut self, other: &RdfCrdt) -> Result<(), OxirsError> {
638 self.triples.merge(&other.triples);
640
641 for (predicate, other_set) in &other.predicate_index {
643 self.predicate_index
644 .entry(predicate.clone())
645 .or_insert_with(|| OrSet::new(self.config.node_id.clone()))
646 .merge(other_set);
647 }
648
649 for (subject, other_set) in &other.subject_index {
651 self.subject_index
652 .entry(subject.clone())
653 .or_insert_with(|| OrSet::new(self.config.node_id.clone()))
654 .merge(other_set);
655 }
656
657 let mut stats = self.stats.write().await;
659 stats.merge_ops += 1;
660 stats.triple_count = self.triples.elements().len();
661
662 Ok(())
663 }
664
665 pub fn get_delta(&self) -> RdfCrdtDelta {
667 RdfCrdtDelta {
668 triples_delta: self.triples.delta(),
669 predicate_deltas: self
670 .predicate_index
671 .iter()
672 .filter_map(|(p, set)| set.delta().map(|d| (p.clone(), d)))
673 .collect(),
674 subject_deltas: self
675 .subject_index
676 .iter()
677 .filter_map(|(s, set)| set.delta().map(|d| (s.clone(), d)))
678 .collect(),
679 }
680 }
681
682 pub async fn apply_delta(&mut self, delta: RdfCrdtDelta) -> Result<(), OxirsError> {
684 if let Some(triples_delta) = delta.triples_delta {
686 self.triples.apply_delta(triples_delta);
687 }
688
689 for (predicate, pred_delta) in delta.predicate_deltas {
691 self.predicate_index
692 .entry(predicate)
693 .or_insert_with(|| OrSet::new(self.config.node_id.clone()))
694 .apply_delta(pred_delta);
695 }
696
697 for (subject, subj_delta) in delta.subject_deltas {
699 self.subject_index
700 .entry(subject)
701 .or_insert_with(|| OrSet::new(self.config.node_id.clone()))
702 .apply_delta(subj_delta);
703 }
704
705 let mut stats = self.stats.write().await;
707 stats.triple_count = self.triples.elements().len();
708
709 Ok(())
710 }
711
712 pub fn reset_delta(&mut self) {
714 self.triples.reset_delta();
715 for set in self.predicate_index.values_mut() {
716 set.reset_delta();
717 }
718 for set in self.subject_index.values_mut() {
719 set.reset_delta();
720 }
721 }
722
723 pub async fn garbage_collect(&mut self) -> Result<GcReport, OxirsError> {
725 let start_tombstones = self.triples.tombstones.len();
726
727 let cutoff = std::time::SystemTime::now()
729 .duration_since(std::time::UNIX_EPOCH)
730 .expect("system clock should be after Unix epoch")
731 .as_secs()
732 - self.config.gc_config.tombstone_ttl_secs;
733
734 self.triples
736 .tombstones
737 .retain(|_, tags| tags.iter().any(|tag| tag.timestamp > cutoff));
738
739 for set in self.predicate_index.values_mut() {
741 set.tombstones
742 .retain(|_, tags| tags.iter().any(|tag| tag.timestamp > cutoff));
743 }
744
745 for set in self.subject_index.values_mut() {
746 set.tombstones
747 .retain(|_, tags| tags.iter().any(|tag| tag.timestamp > cutoff));
748 }
749
750 let removed = start_tombstones - self.triples.tombstones.len();
751
752 Ok(GcReport {
753 tombstones_removed: removed,
754 space_reclaimed: removed * std::mem::size_of::<(Triple, BTreeSet<ElementId>)>(),
755 })
756 }
757
758 pub async fn stats(&self) -> CrdtStatsReport {
760 let stats = self.stats.read().await;
761 CrdtStatsReport {
762 total_ops: stats.total_ops,
763 add_ops: stats.add_ops,
764 remove_ops: stats.remove_ops,
765 merge_ops: stats.merge_ops,
766 triple_count: stats.triple_count,
767 tombstone_count: self.triples.tombstones.len(),
768 }
769 }
770}
771
772#[derive(Debug, Clone, Serialize, Deserialize)]
774pub struct RdfCrdtDelta {
775 pub triples_delta: Option<OrSetDelta<Triple>>,
777 pub predicate_deltas: HashMap<String, OrSetDelta<Triple>>,
779 pub subject_deltas: HashMap<String, OrSetDelta<Triple>>,
781}
782
783#[derive(Debug)]
785pub struct GcReport {
786 pub tombstones_removed: usize,
787 pub space_reclaimed: usize,
788}
789
790#[derive(Debug)]
792pub struct CrdtStatsReport {
793 pub total_ops: u64,
794 pub add_ops: u64,
795 pub remove_ops: u64,
796 pub merge_ops: u64,
797 pub triple_count: usize,
798 pub tombstone_count: usize,
799}
800
801#[cfg(test)]
802mod tests {
803 use super::*;
804 use crate::model::{Literal, NamedNode, Object};
805
806 #[tokio::test]
807 async fn test_grow_set() {
808 let mut set1 = GrowSet::new();
809 let mut set2 = GrowSet::new();
810
811 set1.add(1);
812 set1.add(2);
813 set2.add(2);
814 set2.add(3);
815
816 set1.merge(&set2);
817
818 assert!(set1.contains(&1));
819 assert!(set1.contains(&2));
820 assert!(set1.contains(&3));
821 assert_eq!(set1.elements().len(), 3);
822 }
823
824 #[tokio::test]
825 async fn test_or_set() {
826 let mut set1 = OrSet::new("node1".to_string());
827 let mut set2 = OrSet::new("node2".to_string());
828
829 set1.add(1);
830 set1.add(2);
831 set2.add(2);
832 set2.add(3);
833
834 set1.remove(&2);
836
837 set1.merge(&set2);
839
840 assert!(set1.contains(&1));
841 assert!(set1.contains(&2)); assert!(set1.contains(&3));
843 }
844
845 #[tokio::test]
846 async fn test_rdf_crdt() {
847 let config = CrdtConfig {
848 node_id: "node1".to_string(),
849 crdt_type: CrdtType::RdfCrdt,
850 gc_config: GcConfig::default(),
851 delta_config: DeltaConfig::default(),
852 };
853
854 let mut crdt = RdfCrdt::new(config)
855 .await
856 .expect("async operation should succeed");
857
858 let triple1 = Triple::new(
860 NamedNode::new("http://example.org/s1").expect("valid IRI"),
861 NamedNode::new("http://example.org/p1").expect("valid IRI"),
862 Object::Literal(Literal::new("value1")),
863 );
864
865 let triple2 = Triple::new(
866 NamedNode::new("http://example.org/s1").expect("valid IRI"),
867 NamedNode::new("http://example.org/p2").expect("valid IRI"),
868 Object::Literal(Literal::new("value2")),
869 );
870
871 crdt.add_triple(triple1.clone())
872 .await
873 .expect("async operation should succeed");
874 crdt.add_triple(triple2.clone())
875 .await
876 .expect("async operation should succeed");
877
878 let pattern = TriplePattern::new(
880 Some(crate::model::SubjectPattern::NamedNode(
881 NamedNode::new("http://example.org/s1").expect("valid IRI"),
882 )),
883 None,
884 None,
885 );
886
887 let results = crdt
888 .query(&pattern)
889 .await
890 .expect("async operation should succeed");
891 assert_eq!(results.len(), 2);
892
893 crdt.remove_triple(&triple1)
895 .await
896 .expect("async operation should succeed");
897
898 let results = crdt
899 .query(&pattern)
900 .await
901 .expect("async operation should succeed");
902 assert_eq!(results.len(), 1);
903 assert_eq!(results[0], triple2);
904 }
905
906 #[tokio::test]
912 async fn regression_rdf_crdt_distinguishes_quoted_triple_subjects() {
913 let config = CrdtConfig {
914 node_id: "node1".to_string(),
915 crdt_type: CrdtType::RdfCrdt,
916 gc_config: GcConfig::default(),
917 delta_config: DeltaConfig::default(),
918 };
919 let mut crdt = RdfCrdt::new(config)
920 .await
921 .expect("async operation should succeed");
922
923 let inner1 = Triple::new(
924 NamedNode::new("http://example.org/a1").expect("valid IRI"),
925 NamedNode::new("http://example.org/rel").expect("valid IRI"),
926 NamedNode::new("http://example.org/b1").expect("valid IRI"),
927 );
928 let inner2 = Triple::new(
929 NamedNode::new("http://example.org/a2").expect("valid IRI"),
930 NamedNode::new("http://example.org/rel").expect("valid IRI"),
931 NamedNode::new("http://example.org/b2").expect("valid IRI"),
932 );
933
934 let predicate = NamedNode::new("http://example.org/certainty").expect("valid IRI");
935 let triple1 = Triple::new(
936 crate::model::Subject::QuotedTriple(Box::new(crate::model::QuotedTriple::new(
937 inner1.clone(),
938 ))),
939 predicate.clone(),
940 Object::Literal(Literal::new("high")),
941 );
942 let triple2 = Triple::new(
943 crate::model::Subject::QuotedTriple(Box::new(crate::model::QuotedTriple::new(
944 inner2.clone(),
945 ))),
946 predicate,
947 Object::Literal(Literal::new("low")),
948 );
949
950 crdt.add_triple(triple1.clone())
951 .await
952 .expect("async operation should succeed");
953 crdt.add_triple(triple2.clone())
954 .await
955 .expect("async operation should succeed");
956
957 let algebra_pattern_1 = crate::query::algebra::AlgebraTriplePattern::new(
960 crate::query::algebra::TermPattern::NamedNode(
961 NamedNode::new("http://example.org/a1").expect("valid IRI"),
962 ),
963 crate::query::algebra::TermPattern::NamedNode(
964 NamedNode::new("http://example.org/rel").expect("valid IRI"),
965 ),
966 crate::query::algebra::TermPattern::NamedNode(
967 NamedNode::new("http://example.org/b1").expect("valid IRI"),
968 ),
969 );
970 let pattern1 = TriplePattern::new(
971 Some(crate::model::SubjectPattern::QuotedTriple(Box::new(
972 algebra_pattern_1,
973 ))),
974 None,
975 None,
976 );
977
978 let results = crdt
979 .query(&pattern1)
980 .await
981 .expect("async operation should succeed");
982 assert_eq!(
983 results,
984 vec![triple1],
985 "quoted-triple subject query must return only the structurally matching triple"
986 );
987 }
988
989 #[tokio::test]
990 async fn test_rdf_crdt_merge() {
991 let config1 = CrdtConfig {
992 node_id: "node1".to_string(),
993 crdt_type: CrdtType::RdfCrdt,
994 gc_config: GcConfig::default(),
995 delta_config: DeltaConfig::default(),
996 };
997
998 let config2 = CrdtConfig {
999 node_id: "node2".to_string(),
1000 crdt_type: CrdtType::RdfCrdt,
1001 gc_config: GcConfig::default(),
1002 delta_config: DeltaConfig::default(),
1003 };
1004
1005 let mut crdt1 = RdfCrdt::new(config1)
1006 .await
1007 .expect("async operation should succeed");
1008 let mut crdt2 = RdfCrdt::new(config2)
1009 .await
1010 .expect("async operation should succeed");
1011
1012 let triple1 = Triple::new(
1014 NamedNode::new("http://example.org/s1").expect("valid IRI"),
1015 NamedNode::new("http://example.org/p1").expect("valid IRI"),
1016 Object::Literal(Literal::new("value1")),
1017 );
1018
1019 let triple2 = Triple::new(
1020 NamedNode::new("http://example.org/s2").expect("valid IRI"),
1021 NamedNode::new("http://example.org/p2").expect("valid IRI"),
1022 Object::Literal(Literal::new("value2")),
1023 );
1024
1025 crdt1
1026 .add_triple(triple1.clone())
1027 .await
1028 .expect("async operation should succeed");
1029 crdt2
1030 .add_triple(triple2.clone())
1031 .await
1032 .expect("async operation should succeed");
1033
1034 crdt1
1036 .merge(&crdt2)
1037 .await
1038 .expect("async operation should succeed");
1039
1040 let pattern = TriplePattern::new(None, None, None);
1042 let results = crdt1
1043 .query(&pattern)
1044 .await
1045 .expect("async operation should succeed");
1046 assert_eq!(results.len(), 2);
1047 }
1048}