1use std::collections::BTreeMap;
39
40use super::identity::NodeIdentity;
41use super::ownership::{
42 CatalogVersion, CollectionGroupId, CollectionId, OwnershipEpoch, PlacementAuthorityId,
43 RangeBounds, RangeId, RangeOwnership, ReplicaRole, ShardKeyMode, ShardOwnershipCatalog,
44};
45use super::routing::RoutingHint;
46use super::slot::hash_shard_key_to_range_key;
47use crate::replication::{ReceivedSignal, SignalPlaneMessage};
48
49#[derive(Debug, Clone, Copy, PartialEq, Eq)]
51pub enum TopologyAuthority {
52 RoutingMetadataOnly,
55}
56
57#[derive(Debug, Clone, PartialEq, Eq)]
66pub struct TopologyRange {
67 collection: CollectionId,
68 collection_group: CollectionGroupId,
69 placement_authority: PlacementAuthorityId,
70 range_id: RangeId,
71 shard_key_mode: ShardKeyMode,
72 bounds: RangeBounds,
73 owner: NodeIdentity,
74 replicas: Vec<NodeIdentity>,
75 compressed_archive_replicas: Vec<NodeIdentity>,
76 epoch: OwnershipEpoch,
77 version: CatalogVersion,
78}
79
80impl TopologyRange {
81 fn from_ownership(range: &RangeOwnership) -> Self {
82 Self {
83 collection: range.collection().clone(),
84 collection_group: range.collection_group().clone(),
85 placement_authority: range.placement_authority().clone(),
86 range_id: range.range_id(),
87 shard_key_mode: range.shard_key_mode(),
88 bounds: range.bounds().clone(),
89 owner: range.owner().clone(),
90 replicas: range.hot_mirror_replicas().to_vec(),
91 compressed_archive_replicas: range.compressed_archive_replicas().to_vec(),
92 epoch: range.epoch(),
93 version: range.version(),
94 }
95 }
96
97 pub fn collection(&self) -> &CollectionId {
98 &self.collection
99 }
100
101 pub fn collection_group(&self) -> &CollectionGroupId {
102 &self.collection_group
103 }
104
105 pub fn placement_authority(&self) -> &PlacementAuthorityId {
106 &self.placement_authority
107 }
108
109 pub fn range_id(&self) -> RangeId {
110 self.range_id
111 }
112
113 pub fn shard_key_mode(&self) -> ShardKeyMode {
114 self.shard_key_mode
115 }
116
117 pub fn bounds(&self) -> &RangeBounds {
118 &self.bounds
119 }
120
121 pub fn owner(&self) -> &NodeIdentity {
122 &self.owner
123 }
124
125 pub fn range_owner(&self) -> &NodeIdentity {
126 &self.owner
127 }
128
129 pub fn replicas(&self) -> &[NodeIdentity] {
130 &self.replicas
131 }
132
133 pub fn hot_mirror_replicas(&self) -> &[NodeIdentity] {
134 &self.replicas
135 }
136
137 pub fn compressed_archive_replicas(&self) -> &[NodeIdentity] {
138 &self.compressed_archive_replicas
139 }
140
141 pub fn replica_role_of(&self, node: &NodeIdentity) -> Option<ReplicaRole> {
142 if self.replicas.iter().any(|replica| replica == node) {
143 Some(ReplicaRole::HotMirror)
144 } else if self
145 .compressed_archive_replicas
146 .iter()
147 .any(|replica| replica == node)
148 {
149 Some(ReplicaRole::CompressedArchive)
150 } else {
151 None
152 }
153 }
154
155 pub fn promotion_candidates(&self) -> &[NodeIdentity] {
156 &self.replicas
157 }
158
159 pub fn epoch(&self) -> OwnershipEpoch {
163 self.epoch
164 }
165
166 pub fn version(&self) -> CatalogVersion {
167 self.version
168 }
169
170 fn key(&self) -> (CollectionId, RangeId) {
171 (self.collection.clone(), self.range_id)
172 }
173}
174
175#[derive(Debug, Clone, PartialEq, Eq)]
185pub struct TopologySnapshot {
186 version: CatalogVersion,
187 ranges: Vec<TopologyRange>,
188}
189
190impl TopologySnapshot {
191 pub fn authority(&self) -> TopologyAuthority {
193 TopologyAuthority::RoutingMetadataOnly
194 }
195
196 pub fn version(&self) -> CatalogVersion {
198 self.version
199 }
200
201 pub fn ranges(&self) -> &[TopologyRange] {
203 &self.ranges
204 }
205
206 pub fn range(&self, collection: &CollectionId, range_id: RangeId) -> Option<&TopologyRange> {
208 self.ranges
209 .iter()
210 .find(|r| r.collection() == collection && r.range_id() == range_id)
211 }
212
213 pub fn route(&self, collection: &CollectionId, key: &[u8]) -> Option<&TopologyRange> {
216 self.ranges
217 .iter()
218 .find(|r| r.collection() == collection && r.bounds().contains(key))
219 }
220
221 pub fn route_shard_key(
223 &self,
224 collection: &CollectionId,
225 shard_key: &[u8],
226 ) -> Option<&TopologyRange> {
227 match self.shard_key_mode(collection)? {
228 ShardKeyMode::Ordered => self.route(collection, shard_key),
229 ShardKeyMode::Hash => {
230 let range_key = hash_shard_key_to_range_key(shard_key);
231 self.route(collection, &range_key)
232 }
233 }
234 }
235
236 fn shard_key_mode(&self, collection: &CollectionId) -> Option<ShardKeyMode> {
237 self.ranges
238 .iter()
239 .find(|range| range.collection() == collection)
240 .map(TopologyRange::shard_key_mode)
241 }
242}
243
244impl ShardOwnershipCatalog {
245 pub fn topology_snapshot(&self) -> TopologySnapshot {
253 let ranges: Vec<TopologyRange> =
254 self.entries().map(TopologyRange::from_ownership).collect();
255 let version = ranges
256 .iter()
257 .map(TopologyRange::version)
258 .max()
259 .unwrap_or_else(CatalogVersion::initial);
260 TopologySnapshot { version, ranges }
261 }
262}
263
264#[derive(Debug, Clone, Copy, PartialEq, Eq)]
266pub enum RefreshOutcome {
267 Applied { ranges_changed: usize },
270 Ignored,
273}
274
275impl RefreshOutcome {
276 pub fn was_applied(self) -> bool {
277 matches!(self, RefreshOutcome::Applied { .. })
278 }
279}
280
281#[derive(Debug, Clone, Copy, PartialEq, Eq)]
284pub enum HintOutcome {
285 Corrected,
288 AlreadyCurrent,
291 UnknownRange,
296}
297
298#[derive(Debug, Clone, PartialEq, Eq)]
304pub enum TopologyUpdate {
305 Full(TopologySnapshot),
307 Range(TopologyRange),
309}
310
311#[derive(Debug, Clone, PartialEq, Eq)]
325pub struct ClientTopology {
326 version: CatalogVersion,
327 ranges: BTreeMap<(CollectionId, RangeId), TopologyRange>,
328 needs_refresh: bool,
329}
330
331impl ClientTopology {
332 pub fn from_snapshot(snapshot: TopologySnapshot) -> Self {
334 let mut cache = Self {
335 version: snapshot.version(),
336 ranges: BTreeMap::new(),
337 needs_refresh: false,
338 };
339 for range in snapshot.ranges {
340 cache.ranges.insert(range.key(), range);
341 }
342 cache
343 }
344
345 pub fn version(&self) -> CatalogVersion {
349 self.version
350 }
351
352 pub fn needs_refresh(&self) -> bool {
357 self.needs_refresh
358 }
359
360 pub fn route(&self, collection: &CollectionId, key: &[u8]) -> Option<&TopologyRange> {
362 self.ranges
363 .values()
364 .find(|r| r.collection() == collection && r.bounds().contains(key))
365 }
366
367 pub fn resolve(&self, collection: &CollectionId, key: &[u8]) -> Option<&NodeIdentity> {
369 self.route_shard_key(collection, key)
370 .map(TopologyRange::owner)
371 }
372
373 pub fn route_shard_key(
375 &self,
376 collection: &CollectionId,
377 shard_key: &[u8],
378 ) -> Option<&TopologyRange> {
379 match self.shard_key_mode(collection)? {
380 ShardKeyMode::Ordered => self.route(collection, shard_key),
381 ShardKeyMode::Hash => {
382 let range_key = hash_shard_key_to_range_key(shard_key);
383 self.route(collection, &range_key)
384 }
385 }
386 }
387
388 fn shard_key_mode(&self, collection: &CollectionId) -> Option<ShardKeyMode> {
389 self.ranges
390 .values()
391 .find(|range| range.collection() == collection)
392 .map(TopologyRange::shard_key_mode)
393 }
394
395 pub fn range(&self, collection: &CollectionId, range_id: RangeId) -> Option<&TopologyRange> {
397 self.ranges.get(&(collection.clone(), range_id))
398 }
399
400 pub fn apply_refresh(&mut self, snapshot: TopologySnapshot) -> RefreshOutcome {
409 if !self.ranges.is_empty() && snapshot.version() < self.version {
410 return RefreshOutcome::Ignored;
411 }
412 if self.snapshot_rolls_back_any_range(&snapshot) {
413 return RefreshOutcome::Ignored;
414 }
415 let mut changed = 0usize;
416 let mut next: BTreeMap<(CollectionId, RangeId), TopologyRange> = BTreeMap::new();
417 for range in snapshot.ranges {
418 let key = range.key();
419 if self.ranges.get(&key) != Some(&range) {
420 changed += 1;
421 }
422 next.insert(key, range);
423 }
424 if !self.ranges.is_empty() && snapshot.version <= self.version && changed == 0 {
425 return RefreshOutcome::Ignored;
426 }
427 self.ranges = next;
428 self.version = snapshot.version;
429 self.needs_refresh = false;
430 RefreshOutcome::Applied {
431 ranges_changed: changed,
432 }
433 }
434
435 fn snapshot_rolls_back_any_range(&self, snapshot: &TopologySnapshot) -> bool {
436 snapshot.ranges().iter().any(|incoming| {
437 self.ranges
438 .get(&incoming.key())
439 .is_some_and(|current| incoming.version() < current.version())
440 })
441 }
442
443 pub fn apply_update(&mut self, update: TopologyUpdate) -> RefreshOutcome {
453 match update {
454 TopologyUpdate::Full(snapshot) => self.apply_refresh(snapshot),
455 TopologyUpdate::Range(range) => {
456 let key = range.key();
457 let newer = match self.ranges.get(&key) {
458 Some(current) => range.version() > current.version(),
459 None => true,
460 };
461 if !newer {
462 return RefreshOutcome::Ignored;
463 }
464 if range.version() > self.version {
465 self.version = range.version();
466 }
467 self.ranges.insert(key, range);
468 RefreshOutcome::Applied { ranges_changed: 1 }
469 }
470 }
471 }
472
473 pub fn apply_hint(&mut self, hint: &RoutingHint) -> HintOutcome {
484 let key = (hint.collection().clone(), hint.range_id());
485 match self.ranges.get_mut(&key) {
486 Some(range) => {
487 if hint.version() <= range.version {
488 return HintOutcome::AlreadyCurrent;
489 }
490 range.owner = hint.owner().clone();
491 range.epoch = hint.epoch();
492 range.version = hint.version();
493 self.needs_refresh = true;
494 HintOutcome::Corrected
495 }
496 None => {
497 self.needs_refresh = true;
498 HintOutcome::UnknownRange
499 }
500 }
501 }
502
503 pub fn refresh_on_newer_catalog_signal(
511 &mut self,
512 signals: impl IntoIterator<Item = ReceivedSignal>,
513 mut refresh: impl FnMut() -> TopologySnapshot,
514 ) -> Option<RefreshOutcome> {
515 for signal in signals {
516 let SignalPlaneMessage::CatalogVersionHint(hint) = signal.message else {
517 continue;
518 };
519 if hint.ownership_catalog_version > self.version.value() {
520 return Some(self.apply_refresh(refresh()));
521 }
522 }
523 None
524 }
525}
526
527#[cfg(test)]
528mod tests {
529 use super::*;
530 use crate::cluster::ownership::{PlacementMetadata, RangeBound, ShardKeyMode};
531 use crate::cluster::routing::{RequestOperation, RouteDecision, RoutedRequest, RoutingPolicy};
532
533 fn collection(name: &str) -> CollectionId {
534 CollectionId::new(name).unwrap()
535 }
536
537 fn ident(cn: &str) -> NodeIdentity {
538 NodeIdentity::from_certificate_subject(cn).unwrap()
539 }
540
541 fn full_range(coll: &CollectionId, id: u64, owner: &str, replicas: &[&str]) -> RangeOwnership {
542 RangeOwnership::establish(
543 coll.clone(),
544 RangeId::new(id),
545 ShardKeyMode::Hash,
546 RangeBounds::full(),
547 ident(owner),
548 replicas.iter().map(|r| ident(r)).collect::<Vec<_>>(),
549 PlacementMetadata::with_replication_factor(3),
550 )
551 }
552
553 fn split_range(
554 coll: &CollectionId,
555 id: u64,
556 lower: RangeBound,
557 upper: RangeBound,
558 owner: &str,
559 ) -> RangeOwnership {
560 RangeOwnership::establish(
561 coll.clone(),
562 RangeId::new(id),
563 ShardKeyMode::Ordered,
564 RangeBounds::new(lower, upper).unwrap(),
565 ident(owner),
566 Vec::<NodeIdentity>::new(),
567 PlacementMetadata::with_replication_factor(1),
568 )
569 }
570
571 fn single_hash_slot_bounds(key: &[u8]) -> RangeBounds {
572 let slot = super::super::slot::hash_shard_key_to_slot(key);
573 let lower = RangeBound::key(slot.range_key());
574 let upper = match slot.value().checked_add(1) {
575 Some(next) if next < super::super::slot::PRODUCTION_HASH_SLOT_COUNT => {
576 RangeBound::key(super::super::slot::HashSlot::new(next).unwrap().range_key())
577 }
578 _ => RangeBound::Max,
579 };
580 RangeBounds::new(lower, upper).unwrap()
581 }
582
583 fn hash_slot_range(
584 coll: &CollectionId,
585 id: u64,
586 shard_key: &[u8],
587 owner: &str,
588 ) -> RangeOwnership {
589 RangeOwnership::establish(
590 coll.clone(),
591 RangeId::new(id),
592 ShardKeyMode::Hash,
593 single_hash_slot_bounds(shard_key),
594 ident(owner),
595 Vec::<NodeIdentity>::new(),
596 PlacementMetadata::with_replication_factor(1),
597 )
598 }
599
600 fn catalog_with(ranges: impl IntoIterator<Item = RangeOwnership>) -> ShardOwnershipCatalog {
601 let mut catalog = ShardOwnershipCatalog::new();
602 for range in ranges {
603 catalog.apply_update(range).unwrap();
604 }
605 catalog
606 }
607
608 #[test]
611 fn snapshot_exposes_routing_metadata_for_direct_routing() {
612 let orders = collection("orders");
613 let catalog = catalog_with([full_range(&orders, 1, "CN=node-a", &["CN=node-b"])]);
614
615 let snapshot = catalog.topology_snapshot();
616 assert_eq!(snapshot.version(), CatalogVersion::initial());
617 assert_eq!(snapshot.ranges().len(), 1);
618
619 let range = snapshot
620 .route(&orders, b"any-key")
621 .expect("full range covers all keys");
622 assert_eq!(range.owner(), &ident("CN=node-a"));
623 assert_eq!(range.replicas(), &[ident("CN=node-b")]);
624 assert_eq!(range.epoch(), OwnershipEpoch::initial());
625 assert_eq!(range.range_id(), RangeId::new(1));
626 }
627
628 #[test]
629 fn snapshot_exposes_explicit_replica_roles() {
630 let orders = collection("orders");
631 let catalog = catalog_with([full_range(&orders, 1, "CN=node-a", &["CN=node-b"])
632 .with_collection_group(CollectionGroupId::new("commerce").unwrap())
633 .with_placement_authority(PlacementAuthorityId::new("pa-commerce-1").unwrap())
634 .with_compressed_archive_replicas([ident("CN=node-c")])]);
635
636 let snapshot = catalog.topology_snapshot();
637 assert_eq!(snapshot.authority(), TopologyAuthority::RoutingMetadataOnly);
638 let range = snapshot
639 .route(&orders, b"any-key")
640 .expect("full range covers all keys");
641
642 assert_eq!(range.owner(), &ident("CN=node-a"));
643 assert_eq!(range.range_owner(), &ident("CN=node-a"));
644 assert_eq!(
645 range.collection_group(),
646 &CollectionGroupId::new("commerce").unwrap()
647 );
648 assert_eq!(
649 range.placement_authority(),
650 &PlacementAuthorityId::new("pa-commerce-1").unwrap()
651 );
652 assert_eq!(range.hot_mirror_replicas(), &[ident("CN=node-b")]);
653 assert_eq!(range.compressed_archive_replicas(), &[ident("CN=node-c")]);
654 assert_eq!(range.promotion_candidates(), &[ident("CN=node-b")]);
655 assert_eq!(
656 range.replica_role_of(&ident("CN=node-b")),
657 Some(ReplicaRole::HotMirror)
658 );
659 assert_eq!(
660 range.replica_role_of(&ident("CN=node-c")),
661 Some(ReplicaRole::CompressedArchive)
662 );
663 }
664
665 #[test]
666 fn serving_graph_reflects_ownership_transitions_and_replica_roles() {
667 let orders = collection("orders");
668 let range = full_range(&orders, 1, "CN=node-a", &["CN=node-b"])
669 .with_collection_group(CollectionGroupId::new("commerce").unwrap())
670 .with_placement_authority(PlacementAuthorityId::new("pa-commerce-1").unwrap())
671 .with_compressed_archive_replicas([ident("CN=node-c")]);
672 let mut catalog = catalog_with([range]);
673
674 let current = catalog.range(&orders, RangeId::new(1)).unwrap().clone();
675 catalog
676 .apply_update(current.transfer_to(ident("CN=node-b"), [ident("CN=node-a")]))
677 .unwrap();
678 let current = catalog.range(&orders, RangeId::new(1)).unwrap().clone();
679 catalog
680 .apply_update(current.update_replica_roles([ident("CN=node-a")], [ident("CN=node-d")]))
681 .unwrap();
682
683 let projected = catalog
684 .topology_snapshot()
685 .range(&orders, RangeId::new(1))
686 .unwrap()
687 .clone();
688
689 assert_eq!(projected.range_owner(), &ident("CN=node-b"));
690 assert_eq!(projected.hot_mirror_replicas(), &[ident("CN=node-a")]);
691 assert_eq!(
692 projected.compressed_archive_replicas(),
693 &[ident("CN=node-d")]
694 );
695 assert_eq!(
696 projected.placement_authority(),
697 &PlacementAuthorityId::new("pa-commerce-1").unwrap()
698 );
699 }
700
701 #[test]
704 fn snapshot_routes_keys_to_distinct_owners() {
705 let parts = collection("parts");
706 let catalog = catalog_with([
707 split_range(
708 &parts,
709 1,
710 RangeBound::Min,
711 RangeBound::key(b"m"),
712 "CN=node-a",
713 ),
714 split_range(
715 &parts,
716 2,
717 RangeBound::key(b"m"),
718 RangeBound::Max,
719 "CN=node-b",
720 ),
721 ]);
722 let snapshot = catalog.topology_snapshot();
723
724 assert_eq!(
725 snapshot.route(&parts, b"apple").unwrap().owner(),
726 &ident("CN=node-a")
727 );
728 assert_eq!(
729 snapshot.route(&parts, b"zebra").unwrap().owner(),
730 &ident("CN=node-b")
731 );
732 }
733
734 #[test]
736 fn client_resolves_owner_from_polled_snapshot() {
737 let orders = collection("orders");
738 let catalog = catalog_with([full_range(&orders, 1, "CN=node-a", &[])]);
739 let client = ClientTopology::from_snapshot(catalog.topology_snapshot());
740
741 assert_eq!(client.resolve(&orders, b"k").unwrap(), &ident("CN=node-a"));
742 assert!(!client.needs_refresh());
743 }
744
745 #[test]
746 fn client_resolves_hash_collection_by_shard_key_slot() {
747 let orders = collection("orders");
748 let key = b"tenant:42";
749 let catalog = catalog_with([hash_slot_range(&orders, 1, key, "CN=node-a")]);
750 let client = ClientTopology::from_snapshot(catalog.topology_snapshot());
751
752 let routed = client
753 .route_shard_key(&orders, key)
754 .expect("hash slot range covers the logical shard key");
755 assert_eq!(routed.owner(), &ident("CN=node-a"));
756 assert_eq!(client.resolve(&orders, key).unwrap(), &ident("CN=node-a"));
757 }
758
759 #[test]
762 fn refresh_is_monotonic() {
763 let orders = collection("orders");
764 let mut catalog = catalog_with([full_range(&orders, 1, "CN=node-a", &["CN=node-b"])]);
765 let mut client = ClientTopology::from_snapshot(catalog.topology_snapshot());
766 let v1 = client.version();
767
768 let r = catalog.range(&orders, RangeId::new(1)).unwrap().clone();
770 catalog
771 .apply_update(r.transfer_to(ident("CN=node-b"), [ident("CN=node-a")]))
772 .unwrap();
773 let fresh = catalog.topology_snapshot();
774 assert!(fresh.version() > v1);
775
776 assert_eq!(
777 client.apply_refresh(fresh.clone()),
778 RefreshOutcome::Applied { ranges_changed: 1 }
779 );
780 assert_eq!(client.resolve(&orders, b"k").unwrap(), &ident("CN=node-b"));
781
782 assert_eq!(client.apply_refresh(fresh), RefreshOutcome::Ignored);
784 }
785
786 #[test]
787 fn refresh_applies_same_generation_snapshot_when_another_range_changed() {
788 let parts = collection("parts");
789 let mut catalog = catalog_with([
790 split_range(
791 &parts,
792 1,
793 RangeBound::Min,
794 RangeBound::key(b"m"),
795 "CN=node-a",
796 ),
797 split_range(
798 &parts,
799 2,
800 RangeBound::key(b"m"),
801 RangeBound::Max,
802 "CN=node-b",
803 ),
804 ]);
805 let mut client = ClientTopology::from_snapshot(catalog.topology_snapshot());
806
807 let r1 = catalog.range(&parts, RangeId::new(1)).unwrap().clone();
808 catalog
809 .apply_update(r1.transfer_to(ident("CN=node-c"), Vec::<NodeIdentity>::new()))
810 .unwrap();
811 assert_eq!(
812 client.apply_refresh(catalog.topology_snapshot()),
813 RefreshOutcome::Applied { ranges_changed: 1 }
814 );
815 assert_eq!(
816 client.resolve(&parts, b"apple").unwrap(),
817 &ident("CN=node-c")
818 );
819
820 let r2 = catalog.range(&parts, RangeId::new(2)).unwrap().clone();
821 catalog
822 .apply_update(r2.transfer_to(ident("CN=node-d"), Vec::<NodeIdentity>::new()))
823 .unwrap();
824 let same_generation = catalog.topology_snapshot();
825 assert_eq!(same_generation.version(), client.version());
826
827 assert_eq!(
828 client.apply_refresh(same_generation),
829 RefreshOutcome::Applied { ranges_changed: 1 }
830 );
831 assert_eq!(
832 client.resolve(&parts, b"zebra").unwrap(),
833 &ident("CN=node-d")
834 );
835 }
836
837 #[test]
838 fn equal_generation_refresh_does_not_roll_back_a_newer_range() {
839 let parts = collection("parts");
840 let base = catalog_with([
841 split_range(
842 &parts,
843 1,
844 RangeBound::Min,
845 RangeBound::key(b"m"),
846 "CN=node-a",
847 ),
848 split_range(
849 &parts,
850 2,
851 RangeBound::key(b"m"),
852 RangeBound::Max,
853 "CN=node-b",
854 ),
855 ]);
856 let mut current_catalog = base.clone();
857 let mut stale_fork = base;
858 let mut client = ClientTopology::from_snapshot(current_catalog.topology_snapshot());
859
860 let r1 = current_catalog
861 .range(&parts, RangeId::new(1))
862 .unwrap()
863 .clone();
864 current_catalog
865 .apply_update(r1.transfer_to(ident("CN=node-c"), Vec::<NodeIdentity>::new()))
866 .unwrap();
867 assert!(client
868 .apply_refresh(current_catalog.topology_snapshot())
869 .was_applied());
870
871 let r2 = stale_fork.range(&parts, RangeId::new(2)).unwrap().clone();
872 stale_fork
873 .apply_update(r2.transfer_to(ident("CN=node-d"), Vec::<NodeIdentity>::new()))
874 .unwrap();
875 let fork_snapshot = stale_fork.topology_snapshot();
876 assert_eq!(fork_snapshot.version(), client.version());
877
878 assert_eq!(client.apply_refresh(fork_snapshot), RefreshOutcome::Ignored);
879 assert_eq!(
880 client.resolve(&parts, b"apple").unwrap(),
881 &ident("CN=node-c")
882 );
883 assert_eq!(
884 client.resolve(&parts, b"zebra").unwrap(),
885 &ident("CN=node-b")
886 );
887 }
888
889 #[test]
892 fn redirect_hint_corrects_cache_but_is_not_authoritative() {
893 let orders = collection("orders");
894 let mut catalog = catalog_with([full_range(&orders, 1, "CN=node-a", &["CN=node-b"])]);
895 let mut client = ClientTopology::from_snapshot(catalog.topology_snapshot());
896
897 let r = catalog.range(&orders, RangeId::new(1)).unwrap().clone();
899 catalog
900 .apply_update(r.transfer_to(ident("CN=node-b"), [ident("CN=node-a")]))
901 .unwrap();
902
903 let stale_owner = client.resolve(&orders, b"k").unwrap().clone();
905 assert_eq!(stale_owner, ident("CN=node-a"));
906 let request =
907 RoutedRequest::new(orders.clone(), b"k".to_vec(), RequestOperation::Transaction);
908 let hint = match catalog.plan_route(&stale_owner, &request, &RoutingPolicy::forwarding()) {
909 RouteDecision::Redirect { hint, .. } => hint,
910 other => panic!("expected redirect, got {other:?}"),
911 };
912
913 assert_eq!(client.apply_hint(&hint), HintOutcome::Corrected);
915 assert_eq!(client.resolve(&orders, b"k").unwrap(), &ident("CN=node-b"));
916 assert!(
917 client.needs_refresh(),
918 "a hint is advisory, not authoritative"
919 );
920
921 assert!(client
923 .apply_refresh(catalog.topology_snapshot())
924 .was_applied());
925 assert!(!client.needs_refresh());
926 let range = client.range(&orders, RangeId::new(1)).unwrap();
928 assert_eq!(range.replicas(), &[ident("CN=node-a")]);
929 }
930
931 #[test]
934 fn hint_for_unknown_range_does_not_invent_topology() {
935 let orders = collection("orders");
936 let other = collection("other");
937 let catalog = catalog_with([full_range(&orders, 1, "CN=node-a", &[])]);
938 let mut client = ClientTopology::from_snapshot(catalog.topology_snapshot());
939
940 let foreign = catalog_with([full_range(&other, 9, "CN=node-z", &[])]);
942 let request =
943 RoutedRequest::new(other.clone(), b"k".to_vec(), RequestOperation::Transaction);
944 let hint = foreign
945 .plan_route(&ident("CN=node-b"), &request, &RoutingPolicy::forwarding())
946 .hint()
947 .cloned()
948 .unwrap();
949
950 assert_eq!(client.apply_hint(&hint), HintOutcome::UnknownRange);
951 assert!(
952 client.range(&other, RangeId::new(9)).is_none(),
953 "no phantom range"
954 );
955 assert!(client.needs_refresh());
956 }
957
958 #[test]
961 fn stale_hint_is_ignored_after_authoritative_catch_up() {
962 let orders = collection("orders");
963 let mut catalog = catalog_with([full_range(&orders, 1, "CN=node-a", &["CN=node-b"])]);
964 let request =
965 RoutedRequest::new(orders.clone(), b"k".to_vec(), RequestOperation::Transaction);
966 let early_hint = catalog
968 .plan_route(&ident("CN=node-b"), &request, &RoutingPolicy::forwarding())
969 .hint()
970 .cloned()
971 .unwrap();
972
973 let r = catalog.range(&orders, RangeId::new(1)).unwrap().clone();
975 catalog
976 .apply_update(r.transfer_to(ident("CN=node-b"), [ident("CN=node-a")]))
977 .unwrap();
978 let mut client = ClientTopology::from_snapshot(catalog.topology_snapshot());
979
980 assert_eq!(client.apply_hint(&early_hint), HintOutcome::AlreadyCurrent);
982 assert_eq!(client.resolve(&orders, b"k").unwrap(), &ident("CN=node-b"));
983 assert!(!client.needs_refresh());
984 }
985
986 #[test]
988 fn push_full_snapshot_applies_like_a_poll() {
989 let orders = collection("orders");
990 let mut catalog = catalog_with([full_range(&orders, 1, "CN=node-a", &[])]);
991 let mut client = ClientTopology::from_snapshot(catalog.topology_snapshot());
992
993 let r = catalog.range(&orders, RangeId::new(1)).unwrap().clone();
994 catalog
995 .apply_update(r.transfer_to(ident("CN=node-b"), [ident("CN=node-a")]))
996 .unwrap();
997 let update = TopologyUpdate::Full(catalog.topology_snapshot());
998
999 assert!(client.apply_update(update).was_applied());
1000 assert_eq!(client.resolve(&orders, b"k").unwrap(), &ident("CN=node-b"));
1001 }
1002
1003 #[test]
1005 fn push_range_delta_advances_one_range() {
1006 let parts = collection("parts");
1007 let mut catalog = catalog_with([
1008 split_range(
1009 &parts,
1010 1,
1011 RangeBound::Min,
1012 RangeBound::key(b"m"),
1013 "CN=node-a",
1014 ),
1015 split_range(
1016 &parts,
1017 2,
1018 RangeBound::key(b"m"),
1019 RangeBound::Max,
1020 "CN=node-b",
1021 ),
1022 ]);
1023 let mut client = ClientTopology::from_snapshot(catalog.topology_snapshot());
1024
1025 let r2 = catalog.range(&parts, RangeId::new(2)).unwrap().clone();
1027 catalog
1028 .apply_update(r2.transfer_to(ident("CN=node-c"), Vec::<NodeIdentity>::new()))
1029 .unwrap();
1030 let moved = catalog
1031 .topology_snapshot()
1032 .range(&parts, RangeId::new(2))
1033 .unwrap()
1034 .clone();
1035
1036 assert_eq!(
1037 client.apply_update(TopologyUpdate::Range(moved)),
1038 RefreshOutcome::Applied { ranges_changed: 1 }
1039 );
1040 assert_eq!(
1042 client.resolve(&parts, b"apple").unwrap(),
1043 &ident("CN=node-a")
1044 );
1045 assert_eq!(
1046 client.resolve(&parts, b"zebra").unwrap(),
1047 &ident("CN=node-c")
1048 );
1049 }
1050
1051 #[test]
1055 fn missed_push_still_converges_via_hint_and_poll() {
1056 let orders = collection("orders");
1057 let mut catalog = catalog_with([full_range(&orders, 1, "CN=node-a", &["CN=node-b"])]);
1058 let mut client = ClientTopology::from_snapshot(catalog.topology_snapshot());
1059
1060 let r = catalog.range(&orders, RangeId::new(1)).unwrap().clone();
1063 catalog
1064 .apply_update(r.transfer_to(ident("CN=node-b"), [ident("CN=node-a")]))
1065 .unwrap();
1066 let _dropped_push = TopologyUpdate::Full(catalog.topology_snapshot());
1067
1068 assert_eq!(client.resolve(&orders, b"k").unwrap(), &ident("CN=node-a"));
1070
1071 let request =
1073 RoutedRequest::new(orders.clone(), b"k".to_vec(), RequestOperation::Transaction);
1074 let hint = catalog
1075 .plan_route(&ident("CN=node-a"), &request, &RoutingPolicy::forwarding())
1076 .hint()
1077 .cloned()
1078 .unwrap();
1079 assert_eq!(client.apply_hint(&hint), HintOutcome::Corrected);
1080 assert_eq!(client.resolve(&orders, b"k").unwrap(), &ident("CN=node-b"));
1081 assert!(client.needs_refresh());
1082
1083 assert!(client
1085 .apply_refresh(catalog.topology_snapshot())
1086 .was_applied());
1087 assert!(!client.needs_refresh());
1088 }
1089
1090 #[test]
1093 fn out_of_order_push_keeps_newest() {
1094 let orders = collection("orders");
1095 let mut catalog = catalog_with([full_range(&orders, 1, "CN=node-a", &["CN=node-b"])]);
1096 let mut client = ClientTopology::from_snapshot(catalog.topology_snapshot());
1097
1098 let r1 = catalog.range(&orders, RangeId::new(1)).unwrap().clone();
1100 catalog
1101 .apply_update(r1.transfer_to(ident("CN=node-b"), [ident("CN=node-a")]))
1102 .unwrap();
1103 let push_v2 = catalog
1104 .topology_snapshot()
1105 .range(&orders, RangeId::new(1))
1106 .unwrap()
1107 .clone();
1108
1109 let r2 = catalog.range(&orders, RangeId::new(1)).unwrap().clone();
1111 catalog
1112 .apply_update(r2.transfer_to(ident("CN=node-c"), [ident("CN=node-b")]))
1113 .unwrap();
1114 let push_v3 = catalog
1115 .topology_snapshot()
1116 .range(&orders, RangeId::new(1))
1117 .unwrap()
1118 .clone();
1119
1120 assert!(client
1122 .apply_update(TopologyUpdate::Range(push_v3))
1123 .was_applied());
1124 assert_eq!(client.resolve(&orders, b"k").unwrap(), &ident("CN=node-c"));
1125 assert_eq!(
1127 client.apply_update(TopologyUpdate::Range(push_v2)),
1128 RefreshOutcome::Ignored
1129 );
1130 assert_eq!(client.resolve(&orders, b"k").unwrap(), &ident("CN=node-c"));
1131 }
1132}