1use 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#[derive(Debug, Clone, Serialize, Deserialize)]
35pub struct GeoReplicationConfig {
36 pub region_id: String,
38 pub peers: Vec<PeerRegion>,
40 pub sync_interval_ms: u64,
42 pub max_clock_drift_ms: u64,
44 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 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#[derive(Debug, Clone, Serialize, Deserialize)]
112pub struct PeerRegion {
113 pub region_id: String,
115 pub api_url: String,
117 pub healthy: bool,
119 pub last_sync_ms: u64,
121}
122
123#[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#[derive(Debug, Clone, Serialize, Deserialize)]
134pub struct GeoReplicationStatus {
135 pub region_id: String,
137 pub peers: Vec<PeerStatus>,
139 pub events_sent: u64,
141 pub events_received: u64,
143 pub conflicts_resolved: u64,
145 pub current_hlc: HlcTimestamp,
147 pub version_vectors: std::collections::BTreeMap<String, VersionVector>,
149}
150
151#[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#[derive(Debug, Clone, Serialize, Deserialize)]
163pub struct GeoSyncRequest {
164 pub source_region: String,
166 pub events: Vec<ReplicatedEvent>,
168 pub version_vector: VersionVector,
170}
171
172#[derive(Debug, Clone, Serialize, Deserialize)]
174pub struct GeoSyncResponse {
175 pub accepted: usize,
177 pub skipped: usize,
179 pub version_vector: VersionVector,
181}
182
183pub struct GeoReplicationManager {
188 config: GeoReplicationConfig,
190 hlc: Arc<HybridLogicalClock>,
192 resolver: Arc<CrdtResolver>,
194 peer_health: DashMap<String, PeerHealth>,
196 outbound_buffer: DashMap<String, Vec<ReplicatedEvent>>,
198 events_sent: std::sync::atomic::AtomicU64,
200 events_received: std::sync::atomic::AtomicU64,
202 conflicts_resolved: std::sync::atomic::AtomicU64,
204}
205
206impl GeoReplicationManager {
207 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 pub fn region_id(&self) -> &str {
240 &self.config.region_id
241 }
242
243 pub fn hlc(&self) -> &Arc<HybridLogicalClock> {
245 &self.hlc
246 }
247
248 pub fn resolver(&self) -> &Arc<CrdtResolver> {
250 &self.resolver
251 }
252
253 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 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 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 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 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 self.resolver
318 .merge_version_vector(&request.source_region, &request.version_vector);
319
320 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 pub fn set_peer_health(&self, peer_region: &str, health: PeerHealth) {
335 self.peer_health.insert(peer_region.to_string(), health);
336 }
337
338 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 pub fn peers(&self) -> &[PeerRegion] {
347 &self.config.peers
348 }
349
350 pub fn sync_interval(&self) -> Duration {
352 Duration::from_millis(self.config.sync_interval_ms)
353 }
354
355 pub fn batch_size(&self) -> usize {
357 self.config.batch_size
358 }
359
360 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 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) .map(|p| p.region_id.clone())
394 }
395
396 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 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 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 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 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 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 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 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 assert_eq!(us.resolver.seen_count(), 2);
686 assert_eq!(eu.resolver.seen_count(), 2);
687 }
688}