Skip to main content

oxirs_core/distributed/
crdt.rs

1//! CRDTs for conflict-free replicated RDF
2//!
3//! This module implements Conflict-free Replicated Data Types (CRDTs) optimized
4//! for RDF data, enabling eventual consistency without coordination.
5
6#![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/// CRDT configuration
17#[derive(Debug, Clone)]
18pub struct CrdtConfig {
19    /// Node ID for this replica
20    pub node_id: String,
21    /// CRDT type to use
22    pub crdt_type: CrdtType,
23    /// Garbage collection configuration
24    pub gc_config: GcConfig,
25    /// Delta sync configuration
26    pub delta_config: DeltaConfig,
27}
28
29/// CRDT types available
30#[derive(Debug, Clone)]
31pub enum CrdtType {
32    /// Grow-only set (2P-Set without removals)
33    GSet,
34    /// Two-phase set (add and remove)
35    TwoPhaseSet,
36    /// Add-remove partial order
37    AddRemovePartialOrder,
38    /// Observed-remove set
39    OrSet,
40    /// Last-write-wins element set
41    LwwSet,
42    /// Multi-value register for conflicts
43    MvRegister,
44    /// RDF-specific CRDT
45    RdfCrdt,
46}
47
48/// Garbage collection configuration
49#[derive(Debug, Clone)]
50pub struct GcConfig {
51    /// Enable automatic GC
52    pub auto_gc: bool,
53    /// GC interval in seconds
54    pub interval_secs: u64,
55    /// Maximum tombstone age in seconds
56    pub tombstone_ttl_secs: u64,
57    /// Batch size for GC
58    pub batch_size: usize,
59}
60
61impl Default for GcConfig {
62    fn default() -> Self {
63        GcConfig {
64            auto_gc: true,
65            interval_secs: 3600,           // 1 hour
66            tombstone_ttl_secs: 86400 * 7, // 1 week
67            batch_size: 1000,
68        }
69    }
70}
71
72/// Delta sync configuration
73#[derive(Debug, Clone)]
74pub struct DeltaConfig {
75    /// Enable delta synchronization
76    pub enabled: bool,
77    /// Maximum delta size before full sync
78    pub max_delta_size: usize,
79    /// Delta buffer size
80    pub buffer_size: usize,
81    /// Compression for deltas
82    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
96/// Base trait for CRDTs
97pub trait Crdt: Send + Sync {
98    /// Type of delta for this CRDT
99    type Delta: Send + Sync + Clone + Serialize + for<'de> Deserialize<'de>;
100
101    /// Merge with another CRDT state
102    fn merge(&mut self, other: &Self);
103
104    /// Get delta since last checkpoint
105    fn delta(&self) -> Option<Self::Delta>;
106
107    /// Apply delta
108    fn apply_delta(&mut self, delta: Self::Delta);
109
110    /// Reset delta tracking
111    fn reset_delta(&mut self);
112}
113
114/// Unique ID for CRDT elements
115#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
116pub struct ElementId {
117    /// Timestamp (Lamport clock)
118    pub timestamp: u64,
119    /// Node ID
120    pub node_id: String,
121    /// Random component for uniqueness
122    pub random: u64,
123}
124
125impl ElementId {
126    /// Create new element ID
127    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/// Grow-only set CRDT
140#[derive(Debug, Clone)]
141pub struct GrowSet<T: Clone + Ord + Send + Sync> {
142    /// Elements in the set
143    elements: BTreeSet<T>,
144    /// Delta tracking
145    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    /// Create new grow-only set
156    pub fn new() -> Self {
157        GrowSet {
158            elements: BTreeSet::new(),
159            delta_elements: Some(BTreeSet::new()),
160        }
161    }
162
163    /// Add element
164    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    /// Check if contains element
173    pub fn contains(&self, element: &T) -> bool {
174        self.elements.contains(element)
175    }
176
177    /// Get all elements
178    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/// Two-phase set CRDT
208#[derive(Debug, Clone)]
209pub struct TwoPhaseSet<T: Clone + Ord + Send + Sync> {
210    /// Added elements
211    added: BTreeSet<T>,
212    /// Removed elements (tombstones)
213    removed: BTreeSet<T>,
214    /// Delta tracking
215    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    /// Create new two-phase set
229    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    /// Add element
239    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    /// Remove element
248    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    /// Check if contains element
257    pub fn contains(&self, element: &T) -> bool {
258        self.added.contains(element) && !self.removed.contains(element)
259    }
260
261    /// Get current elements
262    pub fn elements(&self) -> BTreeSet<T> {
263        self.added.difference(&self.removed).cloned().collect()
264    }
265}
266
267/// Observed-Remove Set (OR-Set) CRDT
268#[derive(Debug, Clone)]
269pub struct OrSet<T: Clone + Ord + Send + Sync> {
270    /// Elements with their unique tags
271    elements: BTreeMap<T, BTreeSet<ElementId>>,
272    /// Tombstones for removed elements
273    tombstones: BTreeMap<T, BTreeSet<ElementId>>,
274    /// Node ID for generating tags
275    node_id: String,
276    /// Lamport clock
277    clock: u64,
278    /// Delta tracking
279    delta: Option<OrSetDelta<T>>,
280}
281
282#[derive(Debug, Clone, Serialize, Deserialize)]
283pub struct OrSetDelta<T: Clone + Ord> {
284    /// Added elements with tags
285    added: BTreeMap<T, BTreeSet<ElementId>>,
286    /// Removed tombstones
287    removed: BTreeMap<T, BTreeSet<ElementId>>,
288}
289
290impl<T: Clone + Ord + Send + Sync + Serialize + for<'de> Deserialize<'de>> OrSet<T> {
291    /// Create new OR-Set
292    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    /// Add element
306    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    /// Remove element
325    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    /// Check if contains element
337    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                // Element exists if it has tags not in tombstones
341                !tags.is_subset(tombstone_tags)
342            } else {
343                true
344            }
345        } else {
346            false
347        }
348    }
349
350    /// Get current elements
351    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        // Merge elements
365        for (element, tags) in &other.elements {
366            self.elements
367                .entry(element.clone())
368                .or_default()
369                .extend(tags.iter().cloned());
370        }
371
372        // Merge tombstones
373        for (element, tags) in &other.tombstones {
374            self.tombstones
375                .entry(element.clone())
376                .or_default()
377                .extend(tags.iter().cloned());
378        }
379
380        // Remove elements that are fully tombstoned
381        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        // Update clock
399        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        // Apply added elements
408        for (element, tags) in delta.added {
409            self.elements.entry(element).or_default().extend(tags);
410        }
411
412        // Apply removed tombstones
413        for (element, tags) in delta.removed {
414            self.tombstones
415                .entry(element.clone())
416                .or_default()
417                .extend(tags);
418
419            // Remove element if fully tombstoned
420            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
438/// RDF-specific CRDT optimized for triple stores
439pub struct RdfCrdt {
440    /// Configuration
441    config: CrdtConfig,
442    /// Triple OR-Set for conflict-free triple management
443    triples: OrSet<Triple>,
444    /// Predicate-indexed sets for efficient queries
445    predicate_index: HashMap<String, OrSet<Triple>>,
446    /// Subject-indexed sets
447    subject_index: HashMap<String, OrSet<Triple>>,
448    /// Statistics
449    stats: Arc<RwLock<CrdtStats>>,
450}
451
452/// CRDT statistics
453#[derive(Debug, Default)]
454struct CrdtStats {
455    /// Total operations
456    total_ops: u64,
457    /// Add operations
458    add_ops: u64,
459    /// Remove operations
460    remove_ops: u64,
461    /// Merge operations
462    merge_ops: u64,
463    /// Current triple count
464    triple_count: usize,
465    /// Tombstone count
466    #[allow(dead_code)]
467    tombstone_count: usize,
468}
469
470/// Build the `subject_index` key for a concrete triple subject.
471///
472/// `Subject::as_str()` (via [`crate::model::RdfTerm::as_str`]) returns the
473/// fixed, non-identifying placeholder `"<<quoted-triple>>"` for *every*
474/// quoted-triple subject; keying the index off that constant would collapse
475/// all quoted-triple-subject triples into a single `subject_index` bucket,
476/// so querying by one specific quoted-triple subject pattern would
477/// incorrectly return triples for every *other* quoted-triple subject too.
478/// Serializing the full `<< s p o >>` content instead keeps distinct quoted
479/// triples in distinct buckets; it can never collide with a NamedNode,
480/// BlankNode, or Variable key since none of those start with `"<<"`.
481fn 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    /// Create new RDF CRDT
492    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    /// Add triple
505    pub async fn add_triple(&mut self, triple: Triple) -> Result<(), OxirsError> {
506        // Add to main set
507        self.triples.add(triple.clone());
508
509        // Update predicate index
510        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        // Update subject index
520        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        // Update stats
527        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    /// Remove triple
536    pub async fn remove_triple(&mut self, triple: &Triple) -> Result<(), OxirsError> {
537        // Remove from main set
538        self.triples.remove(triple);
539
540        // Update predicate index
541        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        // Update subject index
550        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        // Update stats
556        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    /// Query triples by pattern
565    pub async fn query(&self, pattern: &TriplePattern) -> Result<Vec<Triple>, OxirsError> {
566        // `subject_index` is keyed by the *serialized content* of concrete
567        // quoted-triple subjects (see `subject_index_key`), not by the
568        // (possibly variable-containing) `SubjectPattern` being queried, so
569        // a `SubjectPattern::QuotedTriple` pattern can never be looked up by
570        // an exact key match the way a NamedNode/BlankNode/Variable pattern
571        // can. Fall back to a full scan for that shape instead of querying
572        // the index with a key that -- by construction -- will never be
573        // present, which would otherwise silently return an empty result
574        // for a pattern that may have real matches.
575        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                // Use both subject and predicate index
589                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                // Use subject index
601                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                // Use predicate index
613                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                // Full scan
625                self.triples
626                    .elements()
627                    .into_iter()
628                    .filter(|t| pattern.matches(t))
629                    .collect()
630            }
631        };
632
633        Ok(results)
634    }
635
636    /// Merge with another RDF CRDT
637    pub async fn merge(&mut self, other: &RdfCrdt) -> Result<(), OxirsError> {
638        // Merge main triple set
639        self.triples.merge(&other.triples);
640
641        // Merge predicate indexes
642        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        // Merge subject indexes
650        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        // Update stats
658        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    /// Get delta for synchronization
666    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    /// Apply delta from another replica
683    pub async fn apply_delta(&mut self, delta: RdfCrdtDelta) -> Result<(), OxirsError> {
684        // Apply main triple delta
685        if let Some(triples_delta) = delta.triples_delta {
686            self.triples.apply_delta(triples_delta);
687        }
688
689        // Apply predicate deltas
690        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        // Apply subject deltas
698        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        // Update stats
706        let mut stats = self.stats.write().await;
707        stats.triple_count = self.triples.elements().len();
708
709        Ok(())
710    }
711
712    /// Reset delta tracking
713    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    /// Garbage collect tombstones
724    pub async fn garbage_collect(&mut self) -> Result<GcReport, OxirsError> {
725        let start_tombstones = self.triples.tombstones.len();
726
727        // Remove old tombstones based on age
728        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        // Filter tombstones by age
735        self.triples
736            .tombstones
737            .retain(|_, tags| tags.iter().any(|tag| tag.timestamp > cutoff));
738
739        // Same for indexes
740        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    /// Get statistics
759    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/// RDF CRDT delta for efficient synchronization
773#[derive(Debug, Clone, Serialize, Deserialize)]
774pub struct RdfCrdtDelta {
775    /// Main triple set delta
776    pub triples_delta: Option<OrSetDelta<Triple>>,
777    /// Predicate index deltas
778    pub predicate_deltas: HashMap<String, OrSetDelta<Triple>>,
779    /// Subject index deltas
780    pub subject_deltas: HashMap<String, OrSetDelta<Triple>>,
781}
782
783/// Garbage collection report
784#[derive(Debug)]
785pub struct GcReport {
786    pub tombstones_removed: usize,
787    pub space_reclaimed: usize,
788}
789
790/// CRDT statistics report
791#[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        // Remove from set1
835        set1.remove(&2);
836
837        // Merge
838        set1.merge(&set2);
839
840        assert!(set1.contains(&1));
841        assert!(set1.contains(&2)); // Still exists due to set2
842        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        // Add triples
859        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        // Query by subject
879        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        // Remove triple
894        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    /// Regression test: `subject_index` must not collapse distinct
907    /// quoted-triple subjects into a single bucket, and querying by a
908    /// specific quoted-triple subject pattern must return exactly the
909    /// triple whose subject structurally matches -- not every
910    /// quoted-triple-subject triple in the store.
911    #[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        // A ground quoted-triple subject pattern matching only `inner1` must
958        // return exactly `triple1`, not both triples.
959        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        // Add different triples to each
1013        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        // Merge
1035        crdt1
1036            .merge(&crdt2)
1037            .await
1038            .expect("async operation should succeed");
1039
1040        // Both triples should be in crdt1
1041        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}