1use std::time::Duration;
40
41use reddb_wire::replication::CausalBookmark;
42use reddb_wire::topology::{Endpoint, ReplicaInfo};
43
44use crate::topology::ClusterMembership;
45
46pub type BookmarkTarget = CausalBookmark;
48
49#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct WriteEndpoint {
53 pub addr: String,
54 pub region: String,
55 pub term: u64,
56}
57
58impl WriteEndpoint {
59 pub fn serves_term(&self, term: u64) -> bool {
63 self.term == term
64 }
65
66 pub fn endpoint(&self) -> Endpoint {
68 Endpoint {
69 addr: self.addr.clone(),
70 region: self.region.clone(),
71 }
72 }
73}
74
75#[derive(Debug, Clone, PartialEq, Eq)]
78pub struct ReadEndpoint {
79 pub addr: String,
80 pub region: String,
81 pub healthy: bool,
82 pub lag_ms: u32,
84 pub frontier_lsn: u64,
86 pub rebootstrapping: bool,
91 pub bookmark_eligible: bool,
95}
96
97impl ReadEndpoint {
98 fn from_replica(info: &ReplicaInfo, bookmark: BookmarkTarget) -> Self {
99 let eligible =
103 info.healthy && !info.rebootstrapping && info.last_applied_lsn >= bookmark.commit_lsn();
104 Self {
105 addr: info.addr.clone(),
106 region: info.region.clone(),
107 healthy: info.healthy,
108 lag_ms: info.lag_ms,
109 frontier_lsn: info.last_applied_lsn,
110 rebootstrapping: info.rebootstrapping,
111 bookmark_eligible: eligible,
112 }
113 }
114
115 fn endpoint(&self) -> Endpoint {
116 Endpoint {
117 addr: self.addr.clone(),
118 region: self.region.clone(),
119 }
120 }
121}
122
123#[derive(Debug, Clone, Copy, PartialEq, Eq)]
125pub enum RouteKind {
126 EligibleTarget,
129 CaughtUpTarget,
132 FallbackReplica,
135 FallbackPrimary,
138}
139
140impl RouteKind {
141 pub fn is_fallback(self) -> bool {
144 matches!(self, Self::FallbackReplica | Self::FallbackPrimary)
145 }
146}
147
148#[derive(Debug, Clone, PartialEq, Eq)]
151pub struct RouteDecision {
152 pub endpoint: Endpoint,
154 pub kind: RouteKind,
156 pub waited: Duration,
158}
159
160pub trait BookmarkWaiter {
171 fn elapsed(&self) -> Duration;
173 fn poll(&mut self, target_addr: &str) -> u64;
176}
177
178#[derive(Debug, Clone, Default)]
180pub struct CausalReadOptions {
181 pub preferred_region: Option<String>,
185 pub deadline: Duration,
188}
189
190impl CausalReadOptions {
191 pub fn with_deadline(deadline: Duration) -> Self {
193 Self {
194 preferred_region: None,
195 deadline,
196 }
197 }
198
199 pub fn prefer_region(mut self, region: impl Into<String>) -> Self {
201 self.preferred_region = Some(region.into());
202 self
203 }
204}
205
206#[derive(Debug, Clone, PartialEq, Eq)]
209pub struct RoutingTable {
210 write: WriteEndpoint,
211 replicas: Vec<ReplicaInfo>,
212 epoch: u64,
213}
214
215impl RoutingTable {
216 pub fn from_membership(membership: ClusterMembership, term: u64) -> Self {
220 let ClusterMembership {
221 primary,
222 replicas,
223 epoch,
224 } = membership;
225 Self {
226 write: WriteEndpoint {
227 addr: primary.addr,
228 region: primary.region,
229 term,
230 },
231 replicas,
232 epoch,
233 }
234 }
235
236 pub fn epoch(&self) -> u64 {
238 self.epoch
239 }
240
241 pub fn write_endpoint(&self) -> &WriteEndpoint {
243 &self.write
244 }
245
246 pub fn read_endpoints(&self, bookmark: BookmarkTarget) -> Vec<ReadEndpoint> {
250 self.replicas
251 .iter()
252 .map(|r| ReadEndpoint::from_replica(r, bookmark))
253 .collect()
254 }
255
256 fn pick_target_index(&self, preferred_region: Option<&str>) -> Option<usize> {
267 if let Some(region) = preferred_region {
268 if let Some(i) = self
269 .replicas
270 .iter()
271 .position(|r| r.healthy && !r.rebootstrapping && r.region == region)
272 {
273 return Some(i);
274 }
275 }
276 self.replicas
277 .iter()
278 .position(|r| r.healthy && !r.rebootstrapping)
279 }
280
281 pub fn route_causal_read(
286 &self,
287 bookmark: BookmarkTarget,
288 opts: &CausalReadOptions,
289 waiter: &mut dyn BookmarkWaiter,
290 ) -> RouteDecision {
291 let target_idx = self.pick_target_index(opts.preferred_region.as_deref());
292
293 if let Some(idx) = target_idx {
294 let target = &self.replicas[idx];
295
296 if target.last_applied_lsn >= bookmark.commit_lsn() {
299 return RouteDecision {
300 endpoint: Endpoint {
301 addr: target.addr.clone(),
302 region: target.region.clone(),
303 },
304 kind: RouteKind::EligibleTarget,
305 waited: Duration::ZERO,
306 };
307 }
308
309 let addr = target.addr.clone();
312 let region = target.region.clone();
313 while waiter.elapsed() < opts.deadline {
314 let frontier = waiter.poll(&addr);
315 if frontier >= bookmark.commit_lsn() {
316 return RouteDecision {
317 endpoint: Endpoint { addr, region },
318 kind: RouteKind::CaughtUpTarget,
319 waited: waiter.elapsed(),
320 };
321 }
322 }
323
324 return self.fall_back(bookmark, Some(idx), waiter.elapsed());
326 }
327
328 self.fall_back(bookmark, None, waiter.elapsed())
330 }
331
332 fn fall_back(
335 &self,
336 bookmark: BookmarkTarget,
337 exclude: Option<usize>,
338 waited: Duration,
339 ) -> RouteDecision {
340 let caught_up = self.replicas.iter().enumerate().find(|(i, r)| {
341 Some(*i) != exclude
342 && r.healthy
343 && !r.rebootstrapping
344 && r.last_applied_lsn >= bookmark.commit_lsn()
345 });
346 match caught_up {
347 Some((_, r)) => RouteDecision {
348 endpoint: ReadEndpoint::from_replica(r, bookmark).endpoint(),
349 kind: RouteKind::FallbackReplica,
350 waited,
351 },
352 None => RouteDecision {
353 endpoint: self.write.endpoint(),
354 kind: RouteKind::FallbackPrimary,
355 waited,
356 },
357 }
358 }
359}
360
361#[cfg(test)]
362mod tests {
363 use super::*;
364 use reddb_wire::topology::Endpoint as WireEndpoint;
365
366 fn primary() -> WireEndpoint {
367 WireEndpoint {
368 addr: "primary:5050".into(),
369 region: "us-east-1".into(),
370 }
371 }
372
373 fn replica(addr: &str, region: &str, healthy: bool, frontier: u64) -> ReplicaInfo {
374 ReplicaInfo {
375 addr: addr.into(),
376 region: region.into(),
377 healthy,
378 lag_ms: if healthy { 5 } else { u32::MAX },
379 last_applied_lsn: frontier,
380 rebootstrapping: false,
381 }
382 }
383
384 fn rebuilding_replica(addr: &str, region: &str, frontier: u64) -> ReplicaInfo {
388 ReplicaInfo {
389 addr: addr.into(),
390 region: region.into(),
391 healthy: true,
392 lag_ms: 5,
393 last_applied_lsn: frontier,
394 rebootstrapping: true,
395 }
396 }
397
398 fn membership(replicas: Vec<ReplicaInfo>) -> ClusterMembership {
399 ClusterMembership {
400 primary: primary(),
401 replicas,
402 epoch: 3,
403 }
404 }
405
406 struct ScriptedWaiter {
411 steps: Vec<u64>,
412 idx: usize,
413 tick: Duration,
414 elapsed: Duration,
415 polled_addrs: Vec<String>,
416 }
417
418 impl ScriptedWaiter {
419 fn new(steps: Vec<u64>, tick: Duration) -> Self {
420 Self {
421 steps,
422 idx: 0,
423 tick,
424 elapsed: Duration::ZERO,
425 polled_addrs: Vec::new(),
426 }
427 }
428 }
429
430 impl BookmarkWaiter for ScriptedWaiter {
431 fn elapsed(&self) -> Duration {
432 self.elapsed
433 }
434 fn poll(&mut self, target_addr: &str) -> u64 {
435 self.polled_addrs.push(target_addr.to_string());
436 self.elapsed += self.tick;
437 let v = self
438 .steps
439 .get(self.idx)
440 .copied()
441 .or_else(|| self.steps.last().copied())
442 .unwrap_or(0);
443 self.idx += 1;
444 v
445 }
446 }
447
448 #[test]
451 fn write_endpoint_is_primary_keyed_by_term() {
452 let table = RoutingTable::from_membership(membership(vec![]), 9);
453 let w = table.write_endpoint();
454 assert_eq!(w.addr, "primary:5050");
455 assert_eq!(w.region, "us-east-1");
456 assert_eq!(w.term, 9);
457 assert!(w.serves_term(9));
458 assert!(!w.serves_term(10));
459 }
460
461 #[test]
464 fn read_endpoints_carry_frontier_and_eligibility() {
465 let table = RoutingTable::from_membership(
466 membership(vec![
467 replica("r-ahead:5050", "us-east-1", true, 200),
468 replica("r-behind:5050", "us-east-1", true, 90),
469 replica("r-down:5050", "us-west-2", false, 500),
470 ]),
471 1,
472 );
473 let bookmark = BookmarkTarget::new(1, 100);
474 let reads = table.read_endpoints(bookmark);
475 assert_eq!(reads.len(), 3);
476
477 assert_eq!(reads[0].frontier_lsn, 200);
479 assert!(reads[0].bookmark_eligible);
480
481 assert_eq!(reads[1].frontier_lsn, 90);
483 assert!(!reads[1].bookmark_eligible);
484
485 assert!(reads[2].frontier_lsn >= 100);
487 assert!(!reads[2].healthy);
488 assert!(!reads[2].bookmark_eligible);
489 }
490
491 #[test]
492 fn eligibility_boundary_is_inclusive_at_commit_lsn() {
493 let table =
494 RoutingTable::from_membership(membership(vec![replica("r:5050", "r1", true, 100)]), 1);
495 let reads = table.read_endpoints(BookmarkTarget::new(1, 100));
497 assert!(reads[0].bookmark_eligible);
498 }
499
500 #[test]
503 fn rebootstrapping_replica_is_never_bookmark_eligible_despite_frontier() {
504 let table = RoutingTable::from_membership(
507 membership(vec![rebuilding_replica("rebuild:5050", "us-east-1", 999)]),
508 1,
509 );
510 let reads = table.read_endpoints(BookmarkTarget::new(1, 100));
511 assert_eq!(reads[0].frontier_lsn, 999);
512 assert!(reads[0].rebootstrapping);
513 assert!(
514 !reads[0].bookmark_eligible,
515 "a rebuilding node must never be bookmark-eligible"
516 );
517 }
518
519 #[test]
520 fn route_skips_rebuilding_node_and_falls_back_to_caught_up_peer() {
521 let table = RoutingTable::from_membership(
525 membership(vec![
526 rebuilding_replica("rebuild:5050", "us-east-1", 999),
527 replica("caught-up:5050", "us-east-1", true, 300),
528 ]),
529 1,
530 );
531 let mut waiter = ScriptedWaiter::new(vec![], Duration::from_millis(10));
532 let decision = table.route_causal_read(
533 BookmarkTarget::new(1, 100),
534 &CausalReadOptions::with_deadline(Duration::from_millis(500)),
535 &mut waiter,
536 );
537 assert_eq!(decision.endpoint.addr, "caught-up:5050");
540 assert!(!decision.kind.is_fallback());
541 assert!(
542 waiter.polled_addrs.iter().all(|a| a != "rebuild:5050"),
543 "must never poll a rebuilding node"
544 );
545 }
546
547 #[test]
548 fn route_falls_back_to_primary_when_every_replica_is_rebuilding() {
549 let table = RoutingTable::from_membership(
552 membership(vec![
553 rebuilding_replica("rebuild-a:5050", "us-east-1", 999),
554 rebuilding_replica("rebuild-b:5050", "us-west-2", 999),
555 ]),
556 4,
557 );
558 let mut waiter = ScriptedWaiter::new(vec![], Duration::from_millis(10));
559 let decision = table.route_causal_read(
560 BookmarkTarget::new(4, 100),
561 &CausalReadOptions::with_deadline(Duration::from_millis(500)),
562 &mut waiter,
563 );
564 assert_eq!(decision.kind, RouteKind::FallbackPrimary);
565 assert_eq!(decision.endpoint.addr, "primary:5050");
566 assert!(
567 waiter.polled_addrs.is_empty(),
568 "no rebuilding node should be polled"
569 );
570 }
571
572 #[test]
573 fn rebuilding_node_excluded_as_fallback_target() {
574 let table = RoutingTable::from_membership(
579 membership(vec![
580 replica("east-lag:5050", "us-east-1", true, 10),
581 rebuilding_replica("west-rebuild:5050", "us-west-2", 999),
582 ]),
583 2,
584 );
585 let mut waiter = ScriptedWaiter::new(vec![10, 20, 30], Duration::from_millis(10));
586 let decision = table.route_causal_read(
587 BookmarkTarget::new(2, 100),
588 &CausalReadOptions::with_deadline(Duration::from_millis(25)).prefer_region("us-east-1"),
589 &mut waiter,
590 );
591 assert_eq!(decision.kind, RouteKind::FallbackPrimary);
592 assert_eq!(decision.endpoint.addr, "primary:5050");
593 }
594
595 #[test]
598 fn route_picks_eligible_target_without_waiting() {
599 let table = RoutingTable::from_membership(
600 membership(vec![replica("r-ok:5050", "us-east-1", true, 150)]),
601 1,
602 );
603 let mut waiter = ScriptedWaiter::new(vec![], Duration::from_millis(10));
604 let decision = table.route_causal_read(
605 BookmarkTarget::new(1, 100),
606 &CausalReadOptions::with_deadline(Duration::from_millis(500)),
607 &mut waiter,
608 );
609 assert_eq!(decision.kind, RouteKind::EligibleTarget);
610 assert_eq!(decision.endpoint.addr, "r-ok:5050");
611 assert_eq!(decision.waited, Duration::ZERO);
612 assert!(waiter.polled_addrs.is_empty(), "must not poll on fast path");
613 }
614
615 #[test]
616 fn route_prefers_region_for_the_target() {
617 let table = RoutingTable::from_membership(
618 membership(vec![
619 replica("east:5050", "us-east-1", true, 150),
620 replica("west:5050", "us-west-2", true, 150),
621 ]),
622 1,
623 );
624 let mut waiter = ScriptedWaiter::new(vec![], Duration::from_millis(10));
625 let decision = table.route_causal_read(
626 BookmarkTarget::new(1, 100),
627 &CausalReadOptions::with_deadline(Duration::from_millis(500))
628 .prefer_region("us-west-2"),
629 &mut waiter,
630 );
631 assert_eq!(decision.endpoint.addr, "west:5050");
632 }
633
634 #[test]
637 fn route_waits_and_routes_to_target_once_it_catches_up() {
638 let table = RoutingTable::from_membership(
639 membership(vec![replica("r-lag:5050", "us-east-1", true, 50)]),
640 1,
641 );
642 let mut waiter = ScriptedWaiter::new(vec![60, 80, 100], Duration::from_millis(10));
645 let decision = table.route_causal_read(
646 BookmarkTarget::new(1, 100),
647 &CausalReadOptions::with_deadline(Duration::from_millis(500)),
648 &mut waiter,
649 );
650 assert_eq!(decision.kind, RouteKind::CaughtUpTarget);
651 assert_eq!(decision.endpoint.addr, "r-lag:5050");
652 assert_eq!(decision.waited, Duration::from_millis(30));
653 assert_eq!(waiter.polled_addrs.len(), 3);
654 }
655
656 #[test]
659 fn route_falls_back_to_caught_up_replica_when_target_stays_behind() {
660 let table = RoutingTable::from_membership(
661 membership(vec![
662 replica("east-lag:5050", "us-east-1", true, 10),
664 replica("west-ok:5050", "us-west-2", true, 300),
666 ]),
667 1,
668 );
669 let mut waiter = ScriptedWaiter::new(vec![10, 20, 30], Duration::from_millis(10));
672 let decision = table.route_causal_read(
673 BookmarkTarget::new(1, 100),
674 &CausalReadOptions::with_deadline(Duration::from_millis(25)).prefer_region("us-east-1"),
675 &mut waiter,
676 );
677 assert_eq!(decision.kind, RouteKind::FallbackReplica);
678 assert_eq!(decision.endpoint.addr, "west-ok:5050");
679 assert!(decision.waited >= Duration::from_millis(25));
680 assert!(waiter.polled_addrs.iter().all(|a| a == "east-lag:5050"));
682 }
683
684 #[test]
687 fn route_falls_back_to_primary_when_no_replica_is_caught_up() {
688 let table = RoutingTable::from_membership(
689 membership(vec![
690 replica("r1:5050", "us-east-1", true, 10),
691 replica("r2:5050", "us-east-1", true, 20),
692 ]),
693 7,
694 );
695 let mut waiter = ScriptedWaiter::new(vec![10, 20], Duration::from_millis(10));
696 let decision = table.route_causal_read(
697 BookmarkTarget::new(7, 100),
698 &CausalReadOptions::with_deadline(Duration::from_millis(15)),
699 &mut waiter,
700 );
701 assert_eq!(decision.kind, RouteKind::FallbackPrimary);
702 assert_eq!(decision.endpoint.addr, "primary:5050");
703 assert!(decision.kind.is_fallback());
704 }
705
706 #[test]
707 fn route_falls_back_to_primary_when_no_replica_is_healthy() {
708 let table = RoutingTable::from_membership(
709 membership(vec![replica("r-down:5050", "us-east-1", false, 500)]),
710 1,
711 );
712 let mut waiter = ScriptedWaiter::new(vec![], Duration::from_millis(10));
713 let decision = table.route_causal_read(
714 BookmarkTarget::new(1, 100),
715 &CausalReadOptions::with_deadline(Duration::from_millis(500)),
716 &mut waiter,
717 );
718 assert_eq!(decision.kind, RouteKind::FallbackPrimary);
720 assert_eq!(decision.endpoint.addr, "primary:5050");
721 assert!(waiter.polled_addrs.is_empty());
722 }
723
724 #[test]
725 fn route_falls_back_to_primary_when_no_replicas_advertised() {
726 let table = RoutingTable::from_membership(membership(vec![]), 1);
727 let mut waiter = ScriptedWaiter::new(vec![], Duration::from_millis(10));
728 let decision = table.route_causal_read(
729 BookmarkTarget::new(1, 100),
730 &CausalReadOptions::with_deadline(Duration::from_millis(500)),
731 &mut waiter,
732 );
733 assert_eq!(decision.kind, RouteKind::FallbackPrimary);
734 assert_eq!(decision.endpoint.addr, "primary:5050");
735 }
736
737 #[test]
740 fn lagging_replica_never_errors_always_resolves_an_endpoint() {
741 let shapes = vec![
744 vec![],
745 vec![replica("a:5050", "r1", true, 0)],
746 vec![replica("a:5050", "r1", false, 0)],
747 vec![
748 replica("a:5050", "r1", true, 1),
749 replica("b:5050", "r2", true, 999),
750 ],
751 ];
752 for shape in shapes {
753 let table = RoutingTable::from_membership(membership(shape), 1);
754 let mut waiter = ScriptedWaiter::new(vec![0], Duration::from_millis(10));
755 let decision = table.route_causal_read(
756 BookmarkTarget::new(1, 100),
757 &CausalReadOptions::with_deadline(Duration::from_millis(20)),
758 &mut waiter,
759 );
760 assert!(!decision.endpoint.addr.is_empty());
761 }
762 }
763}