Skip to main content

nodedb_lite/sync/
clock.rs

1//! Vector clock for sync handshake and delta tracking.
2//!
3//! Tracks per-peer versions so the edge can tell Origin "give me everything
4//! since my last known state" and Origin can respond with exactly the
5//! missing deltas.
6
7use std::collections::HashMap;
8
9/// Vector clock: maps peer IDs to their latest known counter.
10///
11/// On handshake, the edge sends its clock to Origin. Origin compares
12/// against its own clock and sends back deltas for any peer whose
13/// counter is ahead of the edge's view.
14#[derive(Debug, Clone, Default)]
15pub struct VectorClock {
16    /// `peer_id → counter` (Loro version vector entries).
17    entries: HashMap<u64, u64>,
18}
19
20impl VectorClock {
21    pub fn new() -> Self {
22        Self::default()
23    }
24
25    /// Set the counter for a peer.
26    pub fn set(&mut self, peer_id: u64, counter: u64) {
27        self.entries.insert(peer_id, counter);
28    }
29
30    /// Get the counter for a peer (0 if unknown).
31    pub fn get(&self, peer_id: u64) -> u64 {
32        self.entries.get(&peer_id).copied().unwrap_or(0)
33    }
34
35    /// Advance a peer's counter if the new value is greater.
36    pub fn advance(&mut self, peer_id: u64, counter: u64) {
37        let entry = self.entries.entry(peer_id).or_insert(0);
38        if counter > *entry {
39            *entry = counter;
40        }
41    }
42
43    /// Merge another clock into this one (take max per peer).
44    pub fn merge(&mut self, other: &VectorClock) {
45        for (&peer, &counter) in &other.entries {
46            self.advance(peer, counter);
47        }
48    }
49
50    /// Check if this clock dominates another (all entries >= other's).
51    pub fn dominates(&self, other: &VectorClock) -> bool {
52        other
53            .entries
54            .iter()
55            .all(|(&peer, &counter)| self.get(peer) >= counter)
56    }
57
58    /// Compute which peers have advanced beyond our knowledge.
59    ///
60    /// Returns `(peer_id, their_counter, our_counter)` for peers where
61    /// `remote` is ahead of `self`.
62    pub fn diff(&self, remote: &VectorClock) -> Vec<(u64, u64, u64)> {
63        let mut diffs = Vec::new();
64        for (&peer, &remote_counter) in &remote.entries {
65            let local_counter = self.get(peer);
66            if remote_counter > local_counter {
67                diffs.push((peer, remote_counter, local_counter));
68            }
69        }
70        diffs
71    }
72
73    /// Export as a map of `peer_id_hex → counter` for wire format.
74    pub fn to_wire(&self) -> HashMap<String, u64> {
75        self.entries
76            .iter()
77            .map(|(&peer, &counter)| (format!("{peer:016x}"), counter))
78            .collect()
79    }
80
81    /// Import from wire format `peer_id_hex → counter`.
82    pub fn from_wire(wire: &HashMap<String, u64>) -> Self {
83        let entries = wire
84            .iter()
85            .filter_map(|(hex, &counter)| {
86                u64::from_str_radix(hex, 16)
87                    .ok()
88                    .map(|peer| (peer, counter))
89            })
90            .collect();
91        Self { entries }
92    }
93
94    /// Number of peers tracked.
95    pub fn peer_count(&self) -> usize {
96        self.entries.len()
97    }
98
99    /// All tracked peer IDs.
100    pub fn peers(&self) -> Vec<u64> {
101        self.entries.keys().copied().collect()
102    }
103
104    /// Export as HashMap<u64, u64> for direct access.
105    pub fn as_map(&self) -> &HashMap<u64, u64> {
106        &self.entries
107    }
108}
109
110#[cfg(test)]
111mod tests {
112    use super::*;
113
114    #[test]
115    fn new_clock_is_empty() {
116        let c = VectorClock::new();
117        assert_eq!(c.peer_count(), 0);
118        assert_eq!(c.get(1), 0);
119    }
120
121    #[test]
122    fn set_and_get() {
123        let mut c = VectorClock::new();
124        c.set(1, 42);
125        assert_eq!(c.get(1), 42);
126        assert_eq!(c.get(2), 0);
127    }
128
129    #[test]
130    fn advance_only_increases() {
131        let mut c = VectorClock::new();
132        c.set(1, 10);
133        c.advance(1, 5); // Should not decrease.
134        assert_eq!(c.get(1), 10);
135        c.advance(1, 15); // Should increase.
136        assert_eq!(c.get(1), 15);
137    }
138
139    #[test]
140    fn merge_takes_max() {
141        let mut a = VectorClock::new();
142        a.set(1, 10);
143        a.set(2, 20);
144
145        let mut b = VectorClock::new();
146        b.set(1, 15);
147        b.set(3, 30);
148
149        a.merge(&b);
150        assert_eq!(a.get(1), 15); // b was higher.
151        assert_eq!(a.get(2), 20); // only in a.
152        assert_eq!(a.get(3), 30); // only in b.
153    }
154
155    #[test]
156    fn dominates() {
157        let mut a = VectorClock::new();
158        a.set(1, 10);
159        a.set(2, 20);
160
161        let mut b = VectorClock::new();
162        b.set(1, 5);
163        b.set(2, 20);
164
165        assert!(a.dominates(&b));
166        assert!(!b.dominates(&a)); // b.get(1) < a.get(1).
167    }
168
169    #[test]
170    fn diff_finds_ahead_peers() {
171        let mut local = VectorClock::new();
172        local.set(1, 10);
173        local.set(2, 20);
174
175        let mut remote = VectorClock::new();
176        remote.set(1, 15); // ahead
177        remote.set(2, 20); // same
178        remote.set(3, 5); // new peer
179
180        let diffs = local.diff(&remote);
181        assert_eq!(diffs.len(), 2);
182        // peer 1: remote=15, local=10
183        assert!(diffs.iter().any(|&(p, r, l)| p == 1 && r == 15 && l == 10));
184        // peer 3: remote=5, local=0
185        assert!(diffs.iter().any(|&(p, r, l)| p == 3 && r == 5 && l == 0));
186    }
187
188    #[test]
189    fn wire_roundtrip() {
190        let mut c = VectorClock::new();
191        c.set(1, 42);
192        c.set(255, 100);
193
194        let wire = c.to_wire();
195        let restored = VectorClock::from_wire(&wire);
196        assert_eq!(restored.get(1), 42);
197        assert_eq!(restored.get(255), 100);
198    }
199}