Skip to main content

allsource_core/infrastructure/cluster/
geo_replication.rs

1/// Geo-Replication Manager for cross-region event replication.
2///
3/// Coordinates replication between AllSource Core instances in different
4/// geographic regions. Uses HLC for causal ordering and CRDTs for
5/// conflict-free convergence.
6///
7/// # Architecture
8///
9/// ```text
10/// Region US-East          Region EU-West          Region AP-East
11/// ┌─────────────┐        ┌─────────────┐        ┌─────────────┐
12/// │ Core Leader  │◄──────►│ Core Leader  │◄──────►│ Core Leader  │
13/// │ (writes OK)  │  sync  │ (writes OK)  │  sync  │ (writes OK)  │
14/// │ HLC + CRDT   │        │ HLC + CRDT   │        │ HLC + CRDT   │
15/// └─────────────┘        └─────────────┘        └─────────────┘
16/// ```
17///
18/// Each region is an independent leader that accepts writes locally.
19/// Events are replicated asynchronously to peer regions using the
20/// GeoReplicationManager. CRDT resolution ensures convergence.
21///
22/// # Opt-in
23///
24/// Enabled via `ALLSOURCE_GEO_REPLICATION_ENABLED=true` with:
25/// - `ALLSOURCE_REGION_ID`: this region's identifier (e.g., "us-east-1")
26/// - `ALLSOURCE_GEO_PEERS`: comma-separated peer URLs (e.g., "https://eu.core:3900,https://ap.core:3900")
27use super::crdt::{ConflictResolution, CrdtResolver, ReplicatedEvent, VersionVector};
28use super::hlc::{HlcTimestamp, HybridLogicalClock};
29use dashmap::DashMap;
30use serde::{Deserialize, Serialize};
31use std::{sync::Arc, time::Duration};
32
33/// Configuration for geo-replication.
34#[derive(Debug, Clone, Serialize, Deserialize)]
35pub struct GeoReplicationConfig {
36    /// This region's unique identifier.
37    pub region_id: String,
38    /// Peer region endpoints for replication.
39    pub peers: Vec<PeerRegion>,
40    /// Sync interval for pushing events to peers (ms).
41    pub sync_interval_ms: u64,
42    /// Maximum HLC clock drift tolerance (ms).
43    pub max_clock_drift_ms: u64,
44    /// Batch size for replication sync.
45    pub batch_size: usize,
46}
47
48impl Default for GeoReplicationConfig {
49    fn default() -> Self {
50        Self {
51            region_id: "default".to_string(),
52            peers: vec![],
53            sync_interval_ms: 1000,
54            max_clock_drift_ms: 60_000,
55            batch_size: 100,
56        }
57    }
58}
59
60impl GeoReplicationConfig {
61    /// Load configuration from environment variables.
62    pub fn from_env() -> Option<Self> {
63        let enabled = std::env::var("ALLSOURCE_GEO_REPLICATION_ENABLED").is_ok_and(|v| v == "true");
64
65        if !enabled {
66            return None;
67        }
68
69        let region_id =
70            std::env::var("ALLSOURCE_REGION_ID").unwrap_or_else(|_| "default".to_string());
71
72        let peers_str = std::env::var("ALLSOURCE_GEO_PEERS").unwrap_or_default();
73        let peers: Vec<PeerRegion> = peers_str
74            .split(',')
75            .filter(|s| !s.trim().is_empty())
76            .enumerate()
77            .map(|(i, url)| PeerRegion {
78                region_id: format!("peer-{i}"),
79                api_url: url.trim().to_string(),
80                healthy: true,
81                last_sync_ms: 0,
82            })
83            .collect();
84
85        let sync_interval_ms: u64 = std::env::var("ALLSOURCE_GEO_SYNC_INTERVAL_MS")
86            .ok()
87            .and_then(|v| v.parse().ok())
88            .unwrap_or(1000);
89
90        let max_clock_drift_ms: u64 = std::env::var("ALLSOURCE_GEO_MAX_DRIFT_MS")
91            .ok()
92            .and_then(|v| v.parse().ok())
93            .unwrap_or(60_000);
94
95        let batch_size: usize = std::env::var("ALLSOURCE_GEO_BATCH_SIZE")
96            .ok()
97            .and_then(|v| v.parse().ok())
98            .unwrap_or(100);
99
100        Some(Self {
101            region_id,
102            peers,
103            sync_interval_ms,
104            max_clock_drift_ms,
105            batch_size,
106        })
107    }
108}
109
110/// A peer region endpoint.
111#[derive(Debug, Clone, Serialize, Deserialize)]
112pub struct PeerRegion {
113    /// Region identifier.
114    pub region_id: String,
115    /// HTTP API URL for the peer's Core instance.
116    pub api_url: String,
117    /// Whether the peer is currently healthy.
118    pub healthy: bool,
119    /// Last successful sync timestamp (ms since epoch).
120    pub last_sync_ms: u64,
121}
122
123/// Health status of a peer region.
124#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
125#[serde(rename_all = "lowercase")]
126pub enum PeerHealth {
127    Healthy,
128    Degraded,
129    Unreachable,
130}
131
132/// Status of the geo-replication system.
133#[derive(Debug, Clone, Serialize, Deserialize)]
134pub struct GeoReplicationStatus {
135    /// This region's ID.
136    pub region_id: String,
137    /// Peer regions and their health.
138    pub peers: Vec<PeerStatus>,
139    /// Total events replicated out.
140    pub events_sent: u64,
141    /// Total events received from peers.
142    pub events_received: u64,
143    /// Total conflicts resolved (skipped duplicates).
144    pub conflicts_resolved: u64,
145    /// Current HLC state.
146    pub current_hlc: HlcTimestamp,
147    /// Version vectors per region.
148    pub version_vectors: std::collections::BTreeMap<String, VersionVector>,
149}
150
151/// Status of a single peer region.
152#[derive(Debug, Clone, Serialize, Deserialize)]
153pub struct PeerStatus {
154    pub region_id: String,
155    pub api_url: String,
156    pub health: PeerHealth,
157    pub last_sync_ms: u64,
158    pub replication_lag_ms: u64,
159}
160
161/// Replication sync request sent to a peer.
162#[derive(Debug, Clone, Serialize, Deserialize)]
163pub struct GeoSyncRequest {
164    /// Source region ID.
165    pub source_region: String,
166    /// Events to replicate.
167    pub events: Vec<ReplicatedEvent>,
168    /// Source's version vector for this peer.
169    pub version_vector: VersionVector,
170}
171
172/// Replication sync response from a peer.
173#[derive(Debug, Clone, Serialize, Deserialize)]
174pub struct GeoSyncResponse {
175    /// Number of events accepted.
176    pub accepted: usize,
177    /// Number of events skipped (duplicates).
178    pub skipped: usize,
179    /// Peer's current version vector (for bi-directional sync).
180    pub version_vector: VersionVector,
181}
182
183/// Geo-Replication Manager.
184///
185/// Manages cross-region event replication with CRDT conflict resolution
186/// and HLC-based causal ordering.
187pub struct GeoReplicationManager {
188    /// Configuration.
189    config: GeoReplicationConfig,
190    /// Hybrid Logical Clock for this region.
191    hlc: Arc<HybridLogicalClock>,
192    /// CRDT resolver for conflict resolution.
193    resolver: Arc<CrdtResolver>,
194    /// Peer health tracking.
195    peer_health: DashMap<String, PeerHealth>,
196    /// Outbound event buffer (events pending replication to peers).
197    outbound_buffer: DashMap<String, Vec<ReplicatedEvent>>,
198    /// Counter: events sent.
199    events_sent: std::sync::atomic::AtomicU64,
200    /// Counter: events received.
201    events_received: std::sync::atomic::AtomicU64,
202    /// Counter: conflicts resolved.
203    conflicts_resolved: std::sync::atomic::AtomicU64,
204}
205
206impl GeoReplicationManager {
207    /// Create a new geo-replication manager.
208    pub fn new(config: GeoReplicationConfig) -> Self {
209        let node_id: u32 = std::env::var("ALLSOURCE_NODE_ID")
210            .ok()
211            .and_then(|v| v.parse().ok())
212            .unwrap_or(0);
213
214        let hlc = Arc::new(HybridLogicalClock::with_max_drift(
215            node_id,
216            config.max_clock_drift_ms,
217        ));
218
219        let peer_health = DashMap::new();
220        let outbound_buffer = DashMap::new();
221        for peer in &config.peers {
222            peer_health.insert(peer.region_id.clone(), PeerHealth::Healthy);
223            outbound_buffer.insert(peer.region_id.clone(), Vec::new());
224        }
225
226        Self {
227            config,
228            hlc,
229            resolver: Arc::new(CrdtResolver::new()),
230            peer_health,
231            outbound_buffer,
232            events_sent: std::sync::atomic::AtomicU64::new(0),
233            events_received: std::sync::atomic::AtomicU64::new(0),
234            conflicts_resolved: std::sync::atomic::AtomicU64::new(0),
235        }
236    }
237
238    /// Get this region's ID.
239    pub fn region_id(&self) -> &str {
240        &self.config.region_id
241    }
242
243    /// Get the HLC instance.
244    pub fn hlc(&self) -> &Arc<HybridLogicalClock> {
245        &self.hlc
246    }
247
248    /// Get the CRDT resolver.
249    pub fn resolver(&self) -> &Arc<CrdtResolver> {
250        &self.resolver
251    }
252
253    /// Stamp a new local event with an HLC timestamp.
254    pub fn stamp_event(&self, event_id: &str, event_data: serde_json::Value) -> ReplicatedEvent {
255        let ts = self.hlc.now();
256        ReplicatedEvent {
257            event_id: event_id.to_string(),
258            hlc_timestamp: ts,
259            origin_region: self.config.region_id.clone(),
260            event_data,
261        }
262    }
263
264    /// Queue an event for replication to all peers.
265    pub fn queue_for_replication(&self, event: &ReplicatedEvent) {
266        for mut buffer in self.outbound_buffer.iter_mut() {
267            buffer.value_mut().push(event.clone());
268        }
269    }
270
271    /// Drain the outbound buffer for a specific peer.
272    pub fn drain_outbound(&self, peer_region: &str, max_batch: usize) -> Vec<ReplicatedEvent> {
273        if let Some(mut buffer) = self.outbound_buffer.get_mut(peer_region) {
274            let drain_count = max_batch.min(buffer.len());
275            buffer.drain(..drain_count).collect()
276        } else {
277            vec![]
278        }
279    }
280
281    /// Receive events from a peer region.
282    ///
283    /// Applies CRDT conflict resolution and returns the sync response.
284    pub fn receive_sync(&self, request: &GeoSyncRequest) -> GeoSyncResponse {
285        let mut accepted = 0;
286        let mut skipped = 0;
287
288        for event in &request.events {
289            // Update HLC from remote timestamp
290            if let Err(e) = self.hlc.receive(&event.hlc_timestamp) {
291                tracing::warn!(
292                    "HLC drift violation from region {}: {}",
293                    request.source_region,
294                    e,
295                );
296                skipped += 1;
297                self.conflicts_resolved
298                    .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
299                continue;
300            }
301
302            match self.resolver.resolve_and_accept(event) {
303                ConflictResolution::Accept => {
304                    accepted += 1;
305                    self.events_received
306                        .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
307                }
308                ConflictResolution::Skip => {
309                    skipped += 1;
310                    self.conflicts_resolved
311                        .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
312                }
313            }
314        }
315
316        // Merge remote version vector
317        self.resolver
318            .merge_version_vector(&request.source_region, &request.version_vector);
319
320        // Build our version vector for the response
321        let our_vv = self
322            .resolver
323            .version_vector_for(&self.config.region_id)
324            .unwrap_or_default();
325
326        GeoSyncResponse {
327            accepted,
328            skipped,
329            version_vector: our_vv,
330        }
331    }
332
333    /// Update a peer's health status.
334    pub fn set_peer_health(&self, peer_region: &str, health: PeerHealth) {
335        self.peer_health.insert(peer_region.to_string(), health);
336    }
337
338    /// Get a peer's health status.
339    pub fn peer_health(&self, peer_region: &str) -> PeerHealth {
340        self.peer_health
341            .get(peer_region)
342            .map_or(PeerHealth::Unreachable, |h| *h)
343    }
344
345    /// Get all peers.
346    pub fn peers(&self) -> &[PeerRegion] {
347        &self.config.peers
348    }
349
350    /// Get sync interval.
351    pub fn sync_interval(&self) -> Duration {
352        Duration::from_millis(self.config.sync_interval_ms)
353    }
354
355    /// Get batch size.
356    pub fn batch_size(&self) -> usize {
357        self.config.batch_size
358    }
359
360    /// Build a sync request for a specific peer.
361    pub fn build_sync_request(&self, peer_region: &str) -> Option<GeoSyncRequest> {
362        let events = self.drain_outbound(peer_region, self.config.batch_size);
363        if events.is_empty() {
364            return None;
365        }
366
367        self.events_sent
368            .fetch_add(events.len() as u64, std::sync::atomic::Ordering::Relaxed);
369
370        let vv = self
371            .resolver
372            .version_vector_for(&self.config.region_id)
373            .unwrap_or_default();
374
375        Some(GeoSyncRequest {
376            source_region: self.config.region_id.clone(),
377            events,
378            version_vector: vv,
379        })
380    }
381
382    /// Select the best failover region (lowest replication lag, healthy).
383    pub fn select_failover_region(&self) -> Option<String> {
384        self.config
385            .peers
386            .iter()
387            .filter(|p| {
388                self.peer_health
389                    .get(&p.region_id)
390                    .is_some_and(|h| *h == PeerHealth::Healthy)
391            })
392            .max_by_key(|p| p.last_sync_ms) // most recently synced = least data loss
393            .map(|p| p.region_id.clone())
394    }
395
396    /// Get full replication status.
397    pub fn status(&self) -> GeoReplicationStatus {
398        let peers: Vec<PeerStatus> = self
399            .config
400            .peers
401            .iter()
402            .map(|p| {
403                let health = self
404                    .peer_health
405                    .get(&p.region_id)
406                    .map_or(PeerHealth::Unreachable, |h| *h);
407
408                let lag_ms = if p.last_sync_ms > 0 {
409                    let now_ms = std::time::SystemTime::now()
410                        .duration_since(std::time::UNIX_EPOCH)
411                        .unwrap_or_default()
412                        .as_millis() as u64;
413                    now_ms.saturating_sub(p.last_sync_ms)
414                } else {
415                    0
416                };
417
418                PeerStatus {
419                    region_id: p.region_id.clone(),
420                    api_url: p.api_url.clone(),
421                    health,
422                    last_sync_ms: p.last_sync_ms,
423                    replication_lag_ms: lag_ms,
424                }
425            })
426            .collect();
427
428        GeoReplicationStatus {
429            region_id: self.config.region_id.clone(),
430            peers,
431            events_sent: self.events_sent.load(std::sync::atomic::Ordering::Relaxed),
432            events_received: self
433                .events_received
434                .load(std::sync::atomic::Ordering::Relaxed),
435            conflicts_resolved: self
436                .conflicts_resolved
437                .load(std::sync::atomic::Ordering::Relaxed),
438            current_hlc: self.hlc.current(),
439            version_vectors: self.resolver.all_version_vectors(),
440        }
441    }
442}
443
444#[cfg(test)]
445mod tests {
446    use super::*;
447
448    fn test_config(region: &str) -> GeoReplicationConfig {
449        GeoReplicationConfig {
450            region_id: region.to_string(),
451            peers: vec![PeerRegion {
452                region_id: "eu-west".to_string(),
453                api_url: "https://eu.core:3900".to_string(),
454                healthy: true,
455                last_sync_ms: 0,
456            }],
457            sync_interval_ms: 1000,
458            max_clock_drift_ms: 60_000,
459            batch_size: 100,
460        }
461    }
462
463    #[test]
464    fn test_stamp_event() {
465        let mgr = GeoReplicationManager::new(test_config("us-east"));
466        let event = mgr.stamp_event("evt-1", serde_json::json!({"foo": "bar"}));
467
468        assert_eq!(event.event_id, "evt-1");
469        assert_eq!(event.origin_region, "us-east");
470        assert!(event.hlc_timestamp.physical_ms > 0);
471    }
472
473    #[test]
474    fn test_queue_and_drain() {
475        let mgr = GeoReplicationManager::new(test_config("us-east"));
476        let event = mgr.stamp_event("evt-1", serde_json::json!({}));
477
478        mgr.queue_for_replication(&event);
479
480        let batch = mgr.drain_outbound("eu-west", 10);
481        assert_eq!(batch.len(), 1);
482        assert_eq!(batch[0].event_id, "evt-1");
483
484        // Drained — should be empty now
485        let batch2 = mgr.drain_outbound("eu-west", 10);
486        assert!(batch2.is_empty());
487    }
488
489    #[test]
490    fn test_receive_sync_accepts_new_events() {
491        let mgr = GeoReplicationManager::new(test_config("us-east"));
492
493        let request = GeoSyncRequest {
494            source_region: "eu-west".to_string(),
495            events: vec![ReplicatedEvent {
496                event_id: "evt-remote-1".to_string(),
497                hlc_timestamp: HlcTimestamp::new(
498                    std::time::SystemTime::now()
499                        .duration_since(std::time::UNIX_EPOCH)
500                        .unwrap()
501                        .as_millis() as u64,
502                    0,
503                    2,
504                ),
505                origin_region: "eu-west".to_string(),
506                event_data: serde_json::json!({"source": "eu"}),
507            }],
508            version_vector: VersionVector::new(),
509        };
510
511        let response = mgr.receive_sync(&request);
512        assert_eq!(response.accepted, 1);
513        assert_eq!(response.skipped, 0);
514    }
515
516    #[test]
517    fn test_receive_sync_skips_duplicates() {
518        let mgr = GeoReplicationManager::new(test_config("us-east"));
519        let now_ms = std::time::SystemTime::now()
520            .duration_since(std::time::UNIX_EPOCH)
521            .unwrap()
522            .as_millis() as u64;
523
524        let event = ReplicatedEvent {
525            event_id: "evt-dup".to_string(),
526            hlc_timestamp: HlcTimestamp::new(now_ms, 0, 2),
527            origin_region: "eu-west".to_string(),
528            event_data: serde_json::json!({}),
529        };
530
531        let request = GeoSyncRequest {
532            source_region: "eu-west".to_string(),
533            events: vec![event.clone(), event],
534            version_vector: VersionVector::new(),
535        };
536
537        let response = mgr.receive_sync(&request);
538        assert_eq!(response.accepted, 1);
539        assert_eq!(response.skipped, 1);
540    }
541
542    #[test]
543    fn test_build_sync_request() {
544        let mgr = GeoReplicationManager::new(test_config("us-east"));
545        let event = mgr.stamp_event("evt-1", serde_json::json!({}));
546        mgr.queue_for_replication(&event);
547
548        let req = mgr.build_sync_request("eu-west");
549        assert!(req.is_some());
550        let req = req.unwrap();
551        assert_eq!(req.source_region, "us-east");
552        assert_eq!(req.events.len(), 1);
553    }
554
555    #[test]
556    fn test_build_sync_request_empty() {
557        let mgr = GeoReplicationManager::new(test_config("us-east"));
558        let req = mgr.build_sync_request("eu-west");
559        assert!(req.is_none());
560    }
561
562    #[test]
563    fn test_peer_health_tracking() {
564        let mgr = GeoReplicationManager::new(test_config("us-east"));
565        assert_eq!(mgr.peer_health("eu-west"), PeerHealth::Healthy);
566
567        mgr.set_peer_health("eu-west", PeerHealth::Degraded);
568        assert_eq!(mgr.peer_health("eu-west"), PeerHealth::Degraded);
569
570        mgr.set_peer_health("eu-west", PeerHealth::Unreachable);
571        assert_eq!(mgr.peer_health("eu-west"), PeerHealth::Unreachable);
572    }
573
574    #[test]
575    fn test_select_failover_region() {
576        let config = GeoReplicationConfig {
577            region_id: "us-east".to_string(),
578            peers: vec![
579                PeerRegion {
580                    region_id: "eu-west".to_string(),
581                    api_url: "https://eu.core:3900".to_string(),
582                    healthy: true,
583                    last_sync_ms: 100,
584                },
585                PeerRegion {
586                    region_id: "ap-east".to_string(),
587                    api_url: "https://ap.core:3900".to_string(),
588                    healthy: true,
589                    last_sync_ms: 200,
590                },
591            ],
592            ..Default::default()
593        };
594        let mgr = GeoReplicationManager::new(config);
595
596        // ap-east synced more recently → preferred failover
597        let failover = mgr.select_failover_region();
598        assert_eq!(failover, Some("ap-east".to_string()));
599    }
600
601    #[test]
602    fn test_select_failover_skips_unhealthy() {
603        let config = GeoReplicationConfig {
604            region_id: "us-east".to_string(),
605            peers: vec![
606                PeerRegion {
607                    region_id: "eu-west".to_string(),
608                    api_url: "https://eu.core:3900".to_string(),
609                    healthy: true,
610                    last_sync_ms: 200,
611                },
612                PeerRegion {
613                    region_id: "ap-east".to_string(),
614                    api_url: "https://ap.core:3900".to_string(),
615                    healthy: true,
616                    last_sync_ms: 300,
617                },
618            ],
619            ..Default::default()
620        };
621        let mgr = GeoReplicationManager::new(config);
622        mgr.set_peer_health("ap-east", PeerHealth::Unreachable);
623
624        let failover = mgr.select_failover_region();
625        assert_eq!(failover, Some("eu-west".to_string()));
626    }
627
628    #[test]
629    fn test_status() {
630        let mgr = GeoReplicationManager::new(test_config("us-east"));
631        let status = mgr.status();
632
633        assert_eq!(status.region_id, "us-east");
634        assert_eq!(status.peers.len(), 1);
635        assert_eq!(status.events_sent, 0);
636        assert_eq!(status.events_received, 0);
637        assert_eq!(status.conflicts_resolved, 0);
638    }
639
640    #[test]
641    fn test_two_region_convergence() {
642        // Simulate two regions exchanging events
643        let us = GeoReplicationManager::new(GeoReplicationConfig {
644            region_id: "us-east".to_string(),
645            peers: vec![PeerRegion {
646                region_id: "eu-west".to_string(),
647                api_url: "http://eu:3900".to_string(),
648                healthy: true,
649                last_sync_ms: 0,
650            }],
651            ..Default::default()
652        });
653        let eu = GeoReplicationManager::new(GeoReplicationConfig {
654            region_id: "eu-west".to_string(),
655            peers: vec![PeerRegion {
656                region_id: "us-east".to_string(),
657                api_url: "http://us:3900".to_string(),
658                healthy: true,
659                last_sync_ms: 0,
660            }],
661            ..Default::default()
662        });
663
664        // US writes event 1
665        let evt1 = us.stamp_event("evt-1", serde_json::json!({"from": "us"}));
666        us.resolver.resolve_and_accept(&evt1);
667        us.queue_for_replication(&evt1);
668
669        // EU writes event 2
670        let evt2 = eu.stamp_event("evt-2", serde_json::json!({"from": "eu"}));
671        eu.resolver.resolve_and_accept(&evt2);
672        eu.queue_for_replication(&evt2);
673
674        // US → EU sync
675        let us_req = us.build_sync_request("eu-west").unwrap();
676        let eu_resp = eu.receive_sync(&us_req);
677        assert_eq!(eu_resp.accepted, 1);
678
679        // EU → US sync
680        let eu_req = eu.build_sync_request("us-east").unwrap();
681        let us_resp = us.receive_sync(&eu_req);
682        assert_eq!(us_resp.accepted, 1);
683
684        // Both regions now have both events
685        assert_eq!(us.resolver.seen_count(), 2);
686        assert_eq!(eu.resolver.seen_count(), 2);
687    }
688}