1use std::time::Duration;
40
41use reddb_wire::topology::{Endpoint, ReplicaInfo};
42
43use crate::topology::ClusterMembership;
44
45#[derive(Debug, Clone, Copy, PartialEq, Eq)]
51pub struct BookmarkTarget {
52 pub term: u64,
54 pub commit_lsn: u64,
56}
57
58impl BookmarkTarget {
59 pub fn new(term: u64, commit_lsn: u64) -> Self {
60 Self { term, commit_lsn }
61 }
62}
63
64#[derive(Debug, Clone, PartialEq, Eq)]
67pub struct WriteEndpoint {
68 pub addr: String,
69 pub region: String,
70 pub term: u64,
71}
72
73impl WriteEndpoint {
74 pub fn serves_term(&self, term: u64) -> bool {
78 self.term == term
79 }
80
81 pub fn endpoint(&self) -> Endpoint {
83 Endpoint {
84 addr: self.addr.clone(),
85 region: self.region.clone(),
86 }
87 }
88}
89
90#[derive(Debug, Clone, PartialEq, Eq)]
93pub struct ReadEndpoint {
94 pub addr: String,
95 pub region: String,
96 pub healthy: bool,
97 pub lag_ms: u32,
99 pub frontier_lsn: u64,
101 pub rebootstrapping: bool,
106 pub bookmark_eligible: bool,
110}
111
112impl ReadEndpoint {
113 fn from_replica(info: &ReplicaInfo, bookmark: BookmarkTarget) -> Self {
114 let eligible =
118 info.healthy && !info.rebootstrapping && info.last_applied_lsn >= bookmark.commit_lsn;
119 Self {
120 addr: info.addr.clone(),
121 region: info.region.clone(),
122 healthy: info.healthy,
123 lag_ms: info.lag_ms,
124 frontier_lsn: info.last_applied_lsn,
125 rebootstrapping: info.rebootstrapping,
126 bookmark_eligible: eligible,
127 }
128 }
129
130 fn endpoint(&self) -> Endpoint {
131 Endpoint {
132 addr: self.addr.clone(),
133 region: self.region.clone(),
134 }
135 }
136}
137
138#[derive(Debug, Clone, Copy, PartialEq, Eq)]
140pub enum RouteKind {
141 EligibleTarget,
144 CaughtUpTarget,
147 FallbackReplica,
150 FallbackPrimary,
153}
154
155impl RouteKind {
156 pub fn is_fallback(self) -> bool {
159 matches!(self, Self::FallbackReplica | Self::FallbackPrimary)
160 }
161}
162
163#[derive(Debug, Clone, PartialEq, Eq)]
166pub struct RouteDecision {
167 pub endpoint: Endpoint,
169 pub kind: RouteKind,
171 pub waited: Duration,
173}
174
175pub trait BookmarkWaiter {
186 fn elapsed(&self) -> Duration;
188 fn poll(&mut self, target_addr: &str) -> u64;
191}
192
193#[derive(Debug, Clone, Default)]
195pub struct CausalReadOptions {
196 pub preferred_region: Option<String>,
200 pub deadline: Duration,
203}
204
205impl CausalReadOptions {
206 pub fn with_deadline(deadline: Duration) -> Self {
208 Self {
209 preferred_region: None,
210 deadline,
211 }
212 }
213
214 pub fn prefer_region(mut self, region: impl Into<String>) -> Self {
216 self.preferred_region = Some(region.into());
217 self
218 }
219}
220
221#[derive(Debug, Clone, PartialEq, Eq)]
224pub struct RoutingTable {
225 write: WriteEndpoint,
226 replicas: Vec<ReplicaInfo>,
227 epoch: u64,
228}
229
230impl RoutingTable {
231 pub fn from_membership(membership: ClusterMembership, term: u64) -> Self {
235 let ClusterMembership {
236 primary,
237 replicas,
238 epoch,
239 } = membership;
240 Self {
241 write: WriteEndpoint {
242 addr: primary.addr,
243 region: primary.region,
244 term,
245 },
246 replicas,
247 epoch,
248 }
249 }
250
251 pub fn epoch(&self) -> u64 {
253 self.epoch
254 }
255
256 pub fn write_endpoint(&self) -> &WriteEndpoint {
258 &self.write
259 }
260
261 pub fn read_endpoints(&self, bookmark: BookmarkTarget) -> Vec<ReadEndpoint> {
265 self.replicas
266 .iter()
267 .map(|r| ReadEndpoint::from_replica(r, bookmark))
268 .collect()
269 }
270
271 fn pick_target_index(&self, preferred_region: Option<&str>) -> Option<usize> {
282 if let Some(region) = preferred_region {
283 if let Some(i) = self
284 .replicas
285 .iter()
286 .position(|r| r.healthy && !r.rebootstrapping && r.region == region)
287 {
288 return Some(i);
289 }
290 }
291 self.replicas
292 .iter()
293 .position(|r| r.healthy && !r.rebootstrapping)
294 }
295
296 pub fn route_causal_read(
301 &self,
302 bookmark: BookmarkTarget,
303 opts: &CausalReadOptions,
304 waiter: &mut dyn BookmarkWaiter,
305 ) -> RouteDecision {
306 let target_idx = self.pick_target_index(opts.preferred_region.as_deref());
307
308 if let Some(idx) = target_idx {
309 let target = &self.replicas[idx];
310
311 if target.last_applied_lsn >= bookmark.commit_lsn {
314 return RouteDecision {
315 endpoint: Endpoint {
316 addr: target.addr.clone(),
317 region: target.region.clone(),
318 },
319 kind: RouteKind::EligibleTarget,
320 waited: Duration::ZERO,
321 };
322 }
323
324 let addr = target.addr.clone();
327 let region = target.region.clone();
328 while waiter.elapsed() < opts.deadline {
329 let frontier = waiter.poll(&addr);
330 if frontier >= bookmark.commit_lsn {
331 return RouteDecision {
332 endpoint: Endpoint { addr, region },
333 kind: RouteKind::CaughtUpTarget,
334 waited: waiter.elapsed(),
335 };
336 }
337 }
338
339 return self.fall_back(bookmark, Some(idx), waiter.elapsed());
341 }
342
343 self.fall_back(bookmark, None, waiter.elapsed())
345 }
346
347 fn fall_back(
350 &self,
351 bookmark: BookmarkTarget,
352 exclude: Option<usize>,
353 waited: Duration,
354 ) -> RouteDecision {
355 let caught_up = self.replicas.iter().enumerate().find(|(i, r)| {
356 Some(*i) != exclude
357 && r.healthy
358 && !r.rebootstrapping
359 && r.last_applied_lsn >= bookmark.commit_lsn
360 });
361 match caught_up {
362 Some((_, r)) => RouteDecision {
363 endpoint: ReadEndpoint::from_replica(r, bookmark).endpoint(),
364 kind: RouteKind::FallbackReplica,
365 waited,
366 },
367 None => RouteDecision {
368 endpoint: self.write.endpoint(),
369 kind: RouteKind::FallbackPrimary,
370 waited,
371 },
372 }
373 }
374}
375
376#[cfg(test)]
377mod tests {
378 use super::*;
379 use reddb_wire::topology::Endpoint as WireEndpoint;
380
381 fn primary() -> WireEndpoint {
382 WireEndpoint {
383 addr: "primary:5050".into(),
384 region: "us-east-1".into(),
385 }
386 }
387
388 fn replica(addr: &str, region: &str, healthy: bool, frontier: u64) -> ReplicaInfo {
389 ReplicaInfo {
390 addr: addr.into(),
391 region: region.into(),
392 healthy,
393 lag_ms: if healthy { 5 } else { u32::MAX },
394 last_applied_lsn: frontier,
395 rebootstrapping: false,
396 }
397 }
398
399 fn rebuilding_replica(addr: &str, region: &str, frontier: u64) -> ReplicaInfo {
403 ReplicaInfo {
404 addr: addr.into(),
405 region: region.into(),
406 healthy: true,
407 lag_ms: 5,
408 last_applied_lsn: frontier,
409 rebootstrapping: true,
410 }
411 }
412
413 fn membership(replicas: Vec<ReplicaInfo>) -> ClusterMembership {
414 ClusterMembership {
415 primary: primary(),
416 replicas,
417 epoch: 3,
418 }
419 }
420
421 struct ScriptedWaiter {
426 steps: Vec<u64>,
427 idx: usize,
428 tick: Duration,
429 elapsed: Duration,
430 polled_addrs: Vec<String>,
431 }
432
433 impl ScriptedWaiter {
434 fn new(steps: Vec<u64>, tick: Duration) -> Self {
435 Self {
436 steps,
437 idx: 0,
438 tick,
439 elapsed: Duration::ZERO,
440 polled_addrs: Vec::new(),
441 }
442 }
443 }
444
445 impl BookmarkWaiter for ScriptedWaiter {
446 fn elapsed(&self) -> Duration {
447 self.elapsed
448 }
449 fn poll(&mut self, target_addr: &str) -> u64 {
450 self.polled_addrs.push(target_addr.to_string());
451 self.elapsed += self.tick;
452 let v = self
453 .steps
454 .get(self.idx)
455 .copied()
456 .or_else(|| self.steps.last().copied())
457 .unwrap_or(0);
458 self.idx += 1;
459 v
460 }
461 }
462
463 #[test]
466 fn write_endpoint_is_primary_keyed_by_term() {
467 let table = RoutingTable::from_membership(membership(vec![]), 9);
468 let w = table.write_endpoint();
469 assert_eq!(w.addr, "primary:5050");
470 assert_eq!(w.region, "us-east-1");
471 assert_eq!(w.term, 9);
472 assert!(w.serves_term(9));
473 assert!(!w.serves_term(10));
474 }
475
476 #[test]
479 fn read_endpoints_carry_frontier_and_eligibility() {
480 let table = RoutingTable::from_membership(
481 membership(vec![
482 replica("r-ahead:5050", "us-east-1", true, 200),
483 replica("r-behind:5050", "us-east-1", true, 90),
484 replica("r-down:5050", "us-west-2", false, 500),
485 ]),
486 1,
487 );
488 let bookmark = BookmarkTarget::new(1, 100);
489 let reads = table.read_endpoints(bookmark);
490 assert_eq!(reads.len(), 3);
491
492 assert_eq!(reads[0].frontier_lsn, 200);
494 assert!(reads[0].bookmark_eligible);
495
496 assert_eq!(reads[1].frontier_lsn, 90);
498 assert!(!reads[1].bookmark_eligible);
499
500 assert!(reads[2].frontier_lsn >= 100);
502 assert!(!reads[2].healthy);
503 assert!(!reads[2].bookmark_eligible);
504 }
505
506 #[test]
507 fn eligibility_boundary_is_inclusive_at_commit_lsn() {
508 let table =
509 RoutingTable::from_membership(membership(vec![replica("r:5050", "r1", true, 100)]), 1);
510 let reads = table.read_endpoints(BookmarkTarget::new(1, 100));
512 assert!(reads[0].bookmark_eligible);
513 }
514
515 #[test]
518 fn rebootstrapping_replica_is_never_bookmark_eligible_despite_frontier() {
519 let table = RoutingTable::from_membership(
522 membership(vec![rebuilding_replica("rebuild:5050", "us-east-1", 999)]),
523 1,
524 );
525 let reads = table.read_endpoints(BookmarkTarget::new(1, 100));
526 assert_eq!(reads[0].frontier_lsn, 999);
527 assert!(reads[0].rebootstrapping);
528 assert!(
529 !reads[0].bookmark_eligible,
530 "a rebuilding node must never be bookmark-eligible"
531 );
532 }
533
534 #[test]
535 fn route_skips_rebuilding_node_and_falls_back_to_caught_up_peer() {
536 let table = RoutingTable::from_membership(
540 membership(vec![
541 rebuilding_replica("rebuild:5050", "us-east-1", 999),
542 replica("caught-up:5050", "us-east-1", true, 300),
543 ]),
544 1,
545 );
546 let mut waiter = ScriptedWaiter::new(vec![], Duration::from_millis(10));
547 let decision = table.route_causal_read(
548 BookmarkTarget::new(1, 100),
549 &CausalReadOptions::with_deadline(Duration::from_millis(500)),
550 &mut waiter,
551 );
552 assert_eq!(decision.endpoint.addr, "caught-up:5050");
555 assert!(!decision.kind.is_fallback());
556 assert!(
557 waiter.polled_addrs.iter().all(|a| a != "rebuild:5050"),
558 "must never poll a rebuilding node"
559 );
560 }
561
562 #[test]
563 fn route_falls_back_to_primary_when_every_replica_is_rebuilding() {
564 let table = RoutingTable::from_membership(
567 membership(vec![
568 rebuilding_replica("rebuild-a:5050", "us-east-1", 999),
569 rebuilding_replica("rebuild-b:5050", "us-west-2", 999),
570 ]),
571 4,
572 );
573 let mut waiter = ScriptedWaiter::new(vec![], Duration::from_millis(10));
574 let decision = table.route_causal_read(
575 BookmarkTarget::new(4, 100),
576 &CausalReadOptions::with_deadline(Duration::from_millis(500)),
577 &mut waiter,
578 );
579 assert_eq!(decision.kind, RouteKind::FallbackPrimary);
580 assert_eq!(decision.endpoint.addr, "primary:5050");
581 assert!(
582 waiter.polled_addrs.is_empty(),
583 "no rebuilding node should be polled"
584 );
585 }
586
587 #[test]
588 fn rebuilding_node_excluded_as_fallback_target() {
589 let table = RoutingTable::from_membership(
594 membership(vec![
595 replica("east-lag:5050", "us-east-1", true, 10),
596 rebuilding_replica("west-rebuild:5050", "us-west-2", 999),
597 ]),
598 2,
599 );
600 let mut waiter = ScriptedWaiter::new(vec![10, 20, 30], Duration::from_millis(10));
601 let decision = table.route_causal_read(
602 BookmarkTarget::new(2, 100),
603 &CausalReadOptions::with_deadline(Duration::from_millis(25)).prefer_region("us-east-1"),
604 &mut waiter,
605 );
606 assert_eq!(decision.kind, RouteKind::FallbackPrimary);
607 assert_eq!(decision.endpoint.addr, "primary:5050");
608 }
609
610 #[test]
613 fn route_picks_eligible_target_without_waiting() {
614 let table = RoutingTable::from_membership(
615 membership(vec![replica("r-ok:5050", "us-east-1", true, 150)]),
616 1,
617 );
618 let mut waiter = ScriptedWaiter::new(vec![], Duration::from_millis(10));
619 let decision = table.route_causal_read(
620 BookmarkTarget::new(1, 100),
621 &CausalReadOptions::with_deadline(Duration::from_millis(500)),
622 &mut waiter,
623 );
624 assert_eq!(decision.kind, RouteKind::EligibleTarget);
625 assert_eq!(decision.endpoint.addr, "r-ok:5050");
626 assert_eq!(decision.waited, Duration::ZERO);
627 assert!(waiter.polled_addrs.is_empty(), "must not poll on fast path");
628 }
629
630 #[test]
631 fn route_prefers_region_for_the_target() {
632 let table = RoutingTable::from_membership(
633 membership(vec![
634 replica("east:5050", "us-east-1", true, 150),
635 replica("west:5050", "us-west-2", true, 150),
636 ]),
637 1,
638 );
639 let mut waiter = ScriptedWaiter::new(vec![], Duration::from_millis(10));
640 let decision = table.route_causal_read(
641 BookmarkTarget::new(1, 100),
642 &CausalReadOptions::with_deadline(Duration::from_millis(500))
643 .prefer_region("us-west-2"),
644 &mut waiter,
645 );
646 assert_eq!(decision.endpoint.addr, "west:5050");
647 }
648
649 #[test]
652 fn route_waits_and_routes_to_target_once_it_catches_up() {
653 let table = RoutingTable::from_membership(
654 membership(vec![replica("r-lag:5050", "us-east-1", true, 50)]),
655 1,
656 );
657 let mut waiter = ScriptedWaiter::new(vec![60, 80, 100], Duration::from_millis(10));
660 let decision = table.route_causal_read(
661 BookmarkTarget::new(1, 100),
662 &CausalReadOptions::with_deadline(Duration::from_millis(500)),
663 &mut waiter,
664 );
665 assert_eq!(decision.kind, RouteKind::CaughtUpTarget);
666 assert_eq!(decision.endpoint.addr, "r-lag:5050");
667 assert_eq!(decision.waited, Duration::from_millis(30));
668 assert_eq!(waiter.polled_addrs.len(), 3);
669 }
670
671 #[test]
674 fn route_falls_back_to_caught_up_replica_when_target_stays_behind() {
675 let table = RoutingTable::from_membership(
676 membership(vec![
677 replica("east-lag:5050", "us-east-1", true, 10),
679 replica("west-ok:5050", "us-west-2", true, 300),
681 ]),
682 1,
683 );
684 let mut waiter = ScriptedWaiter::new(vec![10, 20, 30], Duration::from_millis(10));
687 let decision = table.route_causal_read(
688 BookmarkTarget::new(1, 100),
689 &CausalReadOptions::with_deadline(Duration::from_millis(25)).prefer_region("us-east-1"),
690 &mut waiter,
691 );
692 assert_eq!(decision.kind, RouteKind::FallbackReplica);
693 assert_eq!(decision.endpoint.addr, "west-ok:5050");
694 assert!(decision.waited >= Duration::from_millis(25));
695 assert!(waiter.polled_addrs.iter().all(|a| a == "east-lag:5050"));
697 }
698
699 #[test]
702 fn route_falls_back_to_primary_when_no_replica_is_caught_up() {
703 let table = RoutingTable::from_membership(
704 membership(vec![
705 replica("r1:5050", "us-east-1", true, 10),
706 replica("r2:5050", "us-east-1", true, 20),
707 ]),
708 7,
709 );
710 let mut waiter = ScriptedWaiter::new(vec![10, 20], Duration::from_millis(10));
711 let decision = table.route_causal_read(
712 BookmarkTarget::new(7, 100),
713 &CausalReadOptions::with_deadline(Duration::from_millis(15)),
714 &mut waiter,
715 );
716 assert_eq!(decision.kind, RouteKind::FallbackPrimary);
717 assert_eq!(decision.endpoint.addr, "primary:5050");
718 assert!(decision.kind.is_fallback());
719 }
720
721 #[test]
722 fn route_falls_back_to_primary_when_no_replica_is_healthy() {
723 let table = RoutingTable::from_membership(
724 membership(vec![replica("r-down:5050", "us-east-1", false, 500)]),
725 1,
726 );
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);
735 assert_eq!(decision.endpoint.addr, "primary:5050");
736 assert!(waiter.polled_addrs.is_empty());
737 }
738
739 #[test]
740 fn route_falls_back_to_primary_when_no_replicas_advertised() {
741 let table = RoutingTable::from_membership(membership(vec![]), 1);
742 let mut waiter = ScriptedWaiter::new(vec![], Duration::from_millis(10));
743 let decision = table.route_causal_read(
744 BookmarkTarget::new(1, 100),
745 &CausalReadOptions::with_deadline(Duration::from_millis(500)),
746 &mut waiter,
747 );
748 assert_eq!(decision.kind, RouteKind::FallbackPrimary);
749 assert_eq!(decision.endpoint.addr, "primary:5050");
750 }
751
752 #[test]
755 fn lagging_replica_never_errors_always_resolves_an_endpoint() {
756 let shapes = vec![
759 vec![],
760 vec![replica("a:5050", "r1", true, 0)],
761 vec![replica("a:5050", "r1", false, 0)],
762 vec![
763 replica("a:5050", "r1", true, 1),
764 replica("b:5050", "r2", true, 999),
765 ],
766 ];
767 for shape in shapes {
768 let table = RoutingTable::from_membership(membership(shape), 1);
769 let mut waiter = ScriptedWaiter::new(vec![0], Duration::from_millis(10));
770 let decision = table.route_causal_read(
771 BookmarkTarget::new(1, 100),
772 &CausalReadOptions::with_deadline(Duration::from_millis(20)),
773 &mut waiter,
774 );
775 assert!(!decision.endpoint.addr.is_empty());
776 }
777 }
778}