1use core::num::NonZeroU16;
4
5use behavior::{
6 Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, Protocol, User,
7};
8use thiserror::Error;
9
10use crate::DeliveryRoute;
11
12pub enum RouterMessage<Route: DeliveryRoute + Clone + PartialEq, R: RoutingStrategy<Route>> {
18 Add(Route),
20 Remove(Route),
22 Route(<Route::Protocol as Protocol>::Msg),
24 Observe(R::Observation),
26}
27
28#[derive(Error, Clone, PartialEq, Eq)]
33pub enum RouterError<M, O, E> {
34 #[error("routing rejected because no recipient is eligible")]
36 NoEligibleRecipients(M),
37 #[error("routing policy selected index {index} from {members} members")]
40 InvalidSelection {
41 message: M,
43 index: usize,
45 members: usize,
47 },
48 #[error("routing policy rejected an observation")]
50 Policy {
51 observation: O,
53 error: E,
55 },
56}
57
58impl<M: core::fmt::Debug, O, E: core::fmt::Debug> core::fmt::Debug for RouterError<M, O, E> {
59 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
60 match self {
61 Self::NoEligibleRecipients(message) => formatter
62 .debug_tuple("NoEligibleRecipients")
63 .field(message)
64 .finish(),
65 Self::InvalidSelection {
66 message,
67 index,
68 members,
69 } => formatter
70 .debug_struct("InvalidSelection")
71 .field("message", message)
72 .field("index", index)
73 .field("members", members)
74 .finish(),
75 Self::Policy { error, .. } => formatter
76 .debug_struct("Policy")
77 .field("observation", &"<retained>")
78 .field("error", error)
79 .finish(),
80 }
81 }
82}
83
84#[must_use = "a rejected routing observation must be returned to its owner"]
90pub struct RoutingObservationRejection<Observation, Reason> {
91 pub observation: Observation,
93 pub reason: Reason,
95}
96
97pub trait RoutingStrategy<Route: DeliveryRoute + Clone + PartialEq>: Clone {
105 type Observation;
107 type Error;
109
110 fn select(
116 &mut self,
117 members: &[Route],
118 message: &<Route::Protocol as Protocol>::Msg,
119 ) -> Option<usize>;
120
121 fn observe(
129 &mut self,
130 _members: &[Route],
131 observation: Self::Observation,
132 ) -> Result<(), RoutingObservationRejection<Self::Observation, Self::Error>>;
133
134 fn added(&mut self, _recipient: Route) {}
136
137 fn removed(&mut self, _index: usize, _recipient: Route, _remaining: usize) {}
139}
140
141#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
147pub struct RoundRobin {
148 next: usize,
149}
150
151impl<Route: DeliveryRoute + Clone + PartialEq> RoutingStrategy<Route> for RoundRobin {
152 type Observation = Never;
153 type Error = Never;
154
155 fn select(
156 &mut self,
157 members: &[Route],
158 _: &<Route::Protocol as Protocol>::Msg,
159 ) -> Option<usize> {
160 if members.is_empty() {
161 return None;
162 }
163 let selected = self.next % members.len();
164 self.next = (selected + 1) % members.len();
165 Some(selected)
166 }
167
168 fn observe(
169 &mut self,
170 _: &[Route],
171 observation: Never,
172 ) -> Result<(), RoutingObservationRejection<Never, Never>> {
173 match observation {}
174 }
175
176 fn removed(&mut self, index: usize, _: Route, remaining: usize) {
177 if remaining == 0 {
178 self.next = 0;
179 } else {
180 if index < self.next {
181 self.next -= 1;
182 }
183 self.next %= remaining;
184 }
185 }
186}
187
188#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
190pub struct LoadVersion(pub u64);
191
192#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
194pub struct Load(pub u64);
195
196pub struct LoadObservation<Route: DeliveryRoute + Clone + PartialEq> {
198 pub recipient: Route,
200 pub version: LoadVersion,
202 pub load: Load,
204}
205
206#[derive(Debug, Clone, Copy, PartialEq, Eq)]
208pub enum LoadEvidence {
209 Unknown,
211 Observed {
213 version: LoadVersion,
215 load: Load,
217 },
218}
219
220#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)]
226pub enum MemberEvidenceError {
227 #[error("evidence names an unknown routing member")]
229 UnknownRecipient,
230 #[error("routing-member evidence is stale")]
232 Stale,
233 #[error("routing-member evidence conflicts at the committed version")]
235 ConflictingVersion,
236}
237
238impl<Route: DeliveryRoute + Clone + PartialEq> Clone for LoadObservation<Route> {
239 fn clone(&self) -> Self {
240 Self {
241 recipient: self.recipient.clone(),
242 version: self.version,
243 load: self.load,
244 }
245 }
246}
247
248impl<Route: DeliveryRoute + Clone + PartialEq> PartialEq for LoadObservation<Route> {
249 fn eq(&self, other: &Self) -> bool {
250 self.recipient == other.recipient
251 && self.version == other.version
252 && self.load == other.load
253 }
254}
255
256impl<Route: DeliveryRoute + Clone + PartialEq> Eq for LoadObservation<Route> {}
257
258#[derive(Clone)]
268pub struct LeastLoaded {
269 loads: Vec<LoadEvidence>,
270}
271
272impl LeastLoaded {
273 #[must_use]
275 pub const fn new() -> Self {
276 Self { loads: Vec::new() }
277 }
278}
279
280impl Default for LeastLoaded {
281 fn default() -> Self {
282 Self::new()
283 }
284}
285
286impl<Route: DeliveryRoute + Clone + PartialEq> RoutingStrategy<Route> for LeastLoaded {
287 type Observation = LoadObservation<Route>;
288 type Error = MemberEvidenceError;
289
290 fn select(
291 &mut self,
292 members: &[Route],
293 _: &<Route::Protocol as Protocol>::Msg,
294 ) -> Option<usize> {
295 self.loads
296 .iter()
297 .enumerate()
298 .take(members.len())
299 .filter_map(|(index, evidence)| match evidence {
300 LoadEvidence::Unknown => None,
301 LoadEvidence::Observed { load, .. } => Some((index, *load)),
302 })
303 .min_by_key(|(index, load)| (*load, *index))
304 .map(|(index, _)| index)
305 }
306
307 fn observe(
308 &mut self,
309 members: &[Route],
310 observation: Self::Observation,
311 ) -> Result<(), RoutingObservationRejection<Self::Observation, Self::Error>> {
312 let Some(index) = members
313 .iter()
314 .position(|member| member == &observation.recipient)
315 else {
316 return Err(RoutingObservationRejection {
317 observation,
318 reason: MemberEvidenceError::UnknownRecipient,
319 });
320 };
321 let Some(evidence) = self.loads.get_mut(index) else {
322 return Err(RoutingObservationRejection {
323 observation,
324 reason: MemberEvidenceError::UnknownRecipient,
325 });
326 };
327 let LoadEvidence::Observed { version, load } = *evidence else {
328 *evidence = LoadEvidence::Observed {
329 version: observation.version,
330 load: observation.load,
331 };
332 return Ok(());
333 };
334 if observation.version < version {
335 return Err(RoutingObservationRejection {
336 observation,
337 reason: MemberEvidenceError::Stale,
338 });
339 }
340 if observation.version == version {
341 return if observation.load == load {
342 Ok(())
343 } else {
344 Err(RoutingObservationRejection {
345 observation,
346 reason: MemberEvidenceError::ConflictingVersion,
347 })
348 };
349 }
350 *evidence = LoadEvidence::Observed {
351 version: observation.version,
352 load: observation.load,
353 };
354 Ok(())
355 }
356
357 fn added(&mut self, _: Route) {
358 self.loads.push(LoadEvidence::Unknown);
359 }
360
361 fn removed(&mut self, index: usize, _: Route, _: usize) {
362 if index < self.loads.len() {
363 self.loads.remove(index);
364 }
365 }
366}
367
368pub trait RouteKey<K> {
370 fn route_key(&self) -> &K;
372}
373
374#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
376pub struct MemberToken(pub u64);
377
378#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
380pub struct MemberTokenVersion(pub u64);
381
382pub struct MemberTokenObservation<Route: DeliveryRoute + Clone + PartialEq> {
384 pub recipient: Route,
386 pub version: MemberTokenVersion,
388 pub token: MemberToken,
390}
391
392impl<Route: DeliveryRoute + Clone + PartialEq> Clone for MemberTokenObservation<Route> {
393 fn clone(&self) -> Self {
394 Self {
395 recipient: self.recipient.clone(),
396 version: self.version,
397 token: self.token,
398 }
399 }
400}
401
402impl<Route: DeliveryRoute + Clone + PartialEq> PartialEq for MemberTokenObservation<Route> {
403 fn eq(&self, other: &Self) -> bool {
404 self.recipient == other.recipient
405 && self.version == other.version
406 && self.token == other.token
407 }
408}
409
410impl<Route: DeliveryRoute + Clone + PartialEq> Eq for MemberTokenObservation<Route> {}
411
412#[derive(Debug, Clone, Copy, PartialEq, Eq)]
414pub enum MemberTokenEvidence {
415 Unknown,
417 Observed {
419 version: MemberTokenVersion,
421 token: MemberToken,
423 },
424}
425
426#[derive(Clone)]
427struct HashMembership {
428 evidence: Vec<MemberTokenEvidence>,
429}
430
431impl HashMembership {
432 const fn new() -> Self {
433 Self {
434 evidence: Vec::new(),
435 }
436 }
437
438 fn added(&mut self) {
439 self.evidence.push(MemberTokenEvidence::Unknown);
440 }
441
442 fn removed(&mut self, index: usize) {
443 if index < self.evidence.len() {
444 self.evidence.remove(index);
445 }
446 }
447
448 fn observe<Route: DeliveryRoute + Clone + PartialEq>(
449 &mut self,
450 recipients: &[Route],
451 observation: MemberTokenObservation<Route>,
452 ) -> Result<(), RoutingObservationRejection<MemberTokenObservation<Route>, MemberEvidenceError>>
453 {
454 let Some(index) = recipients
455 .iter()
456 .position(|member| member == &observation.recipient)
457 else {
458 return Err(RoutingObservationRejection {
459 observation,
460 reason: MemberEvidenceError::UnknownRecipient,
461 });
462 };
463 let Some(evidence) = self.evidence.get_mut(index) else {
464 return Err(RoutingObservationRejection {
465 observation,
466 reason: MemberEvidenceError::UnknownRecipient,
467 });
468 };
469 let MemberTokenEvidence::Observed { version, token } = *evidence else {
470 *evidence = MemberTokenEvidence::Observed {
471 version: observation.version,
472 token: observation.token,
473 };
474 return Ok(());
475 };
476 if observation.version < version {
477 return Err(RoutingObservationRejection {
478 observation,
479 reason: MemberEvidenceError::Stale,
480 });
481 }
482 if observation.version == version {
483 return if observation.token == token {
484 Ok(())
485 } else {
486 Err(RoutingObservationRejection {
487 observation,
488 reason: MemberEvidenceError::ConflictingVersion,
489 })
490 };
491 }
492 *evidence = MemberTokenEvidence::Observed {
493 version: observation.version,
494 token: observation.token,
495 };
496 Ok(())
497 }
498
499 fn tokens(&self, members: usize) -> impl Iterator<Item = (usize, MemberToken)> + '_ {
500 self.evidence.iter().enumerate().take(members).filter_map(
501 |(index, evidence)| match evidence {
502 MemberTokenEvidence::Unknown => None,
503 MemberTokenEvidence::Observed { token, .. } => Some((index, *token)),
504 },
505 )
506 }
507}
508
509fn mixed_hash(left: u64, right: u64) -> u64 {
510 let mut value = left ^ right.rotate_left(32) ^ 0x9E37_79B9_7F4A_7C15;
511 value = (value ^ (value >> 30)).wrapping_mul(0xBF58_476D_1CE4_E5B9);
512 value = (value ^ (value >> 27)).wrapping_mul(0x94D0_49BB_1331_11EB);
513 value ^ (value >> 31)
514}
515
516pub struct ConsistentHash<K> {
528 membership: HashMembership,
529 replicas: NonZeroU16,
530 hash_key: fn(&K) -> u64,
531}
532
533impl<K> Clone for ConsistentHash<K> {
534 fn clone(&self) -> Self {
535 Self {
536 membership: self.membership.clone(),
537 replicas: self.replicas,
538 hash_key: self.hash_key,
539 }
540 }
541}
542
543impl<K> ConsistentHash<K> {
544 #[must_use]
546 pub const fn new(replicas: NonZeroU16, hash_key: fn(&K) -> u64) -> Self {
547 Self {
548 membership: HashMembership::new(),
549 replicas,
550 hash_key,
551 }
552 }
553}
554
555impl<Route, K> RoutingStrategy<Route> for ConsistentHash<K>
556where
557 Route: DeliveryRoute + Clone + PartialEq,
558 <Route::Protocol as Protocol>::Msg: RouteKey<K>,
559{
560 type Observation = MemberTokenObservation<Route>;
561 type Error = MemberEvidenceError;
562
563 fn select(
564 &mut self,
565 members: &[Route],
566 message: &<Route::Protocol as Protocol>::Msg,
567 ) -> Option<usize> {
568 let key = (self.hash_key)(message.route_key());
569 self.membership
570 .tokens(members.len())
571 .flat_map(|(index, token)| {
572 (0..self.replicas.get())
573 .map(move |replica| (mixed_hash(token.0, u64::from(replica)), index))
574 })
575 .min_by_key(|(point, index)| (*point < key, *point, *index))
576 .map(|(_, index)| index)
577 }
578
579 fn observe(
580 &mut self,
581 members: &[Route],
582 observation: Self::Observation,
583 ) -> Result<(), RoutingObservationRejection<Self::Observation, Self::Error>> {
584 self.membership.observe(members, observation)
585 }
586
587 fn added(&mut self, _: Route) {
588 self.membership.added();
589 }
590
591 fn removed(&mut self, index: usize, _: Route, _: usize) {
592 self.membership.removed(index);
593 }
594}
595
596pub struct RendezvousHash<K> {
606 membership: HashMembership,
607 hash_key: fn(&K) -> u64,
608}
609
610impl<K> Clone for RendezvousHash<K> {
611 fn clone(&self) -> Self {
612 Self {
613 membership: self.membership.clone(),
614 hash_key: self.hash_key,
615 }
616 }
617}
618
619impl<K> RendezvousHash<K> {
620 #[must_use]
622 pub const fn new(hash_key: fn(&K) -> u64) -> Self {
623 Self {
624 membership: HashMembership::new(),
625 hash_key,
626 }
627 }
628}
629
630impl<Route, K> RoutingStrategy<Route> for RendezvousHash<K>
631where
632 Route: DeliveryRoute + Clone + PartialEq,
633 <Route::Protocol as Protocol>::Msg: RouteKey<K>,
634{
635 type Observation = MemberTokenObservation<Route>;
636 type Error = MemberEvidenceError;
637
638 fn select(
639 &mut self,
640 members: &[Route],
641 message: &<Route::Protocol as Protocol>::Msg,
642 ) -> Option<usize> {
643 let key = (self.hash_key)(message.route_key());
644 self.membership
645 .tokens(members.len())
646 .map(|(index, token)| (mixed_hash(key, token.0), index))
647 .max_by(|left, right| left.0.cmp(&right.0).then_with(|| right.1.cmp(&left.1)))
648 .map(|(_, index)| index)
649 }
650
651 fn observe(
652 &mut self,
653 members: &[Route],
654 observation: Self::Observation,
655 ) -> Result<(), RoutingObservationRejection<Self::Observation, Self::Error>> {
656 self.membership.observe(members, observation)
657 }
658
659 fn added(&mut self, _: Route) {
660 self.membership.added();
661 }
662
663 fn removed(&mut self, index: usize, _: Route, _: usize) {
664 self.membership.removed(index);
665 }
666}
667
668pub struct Router<
680 A: Address,
681 Route: DeliveryRoute<Protocol: Protocol<Addr = A>> + Clone + PartialEq,
682 R,
683> {
684 recipients: Vec<Route>,
685 strategy: R,
686}
687
688impl<A, Route, R> Router<A, Route, R>
689where
690 A: Address,
691 Route: DeliveryRoute<Protocol: Protocol<Addr = A>> + Clone + PartialEq,
692 R: RoutingStrategy<Route>,
693{
694 #[must_use]
696 pub fn new(recipients: Vec<Route>, strategy: R) -> Self {
697 let mut unique = Vec::with_capacity(recipients.len());
698 for recipient in recipients {
699 if !unique.contains(&recipient) {
700 unique.push(recipient);
701 }
702 }
703 let mut strategy = strategy;
704 for recipient in &unique {
705 strategy.added(recipient.clone());
706 }
707 Self {
708 recipients: unique,
709 strategy,
710 }
711 }
712
713 #[must_use]
715 pub fn recipients(&self) -> &[Route] {
716 &self.recipients
717 }
718
719 #[must_use]
721 pub const fn strategy(&self) -> &R {
722 &self.strategy
723 }
724
725 fn member_index(&self, recipient: &Route) -> Option<usize> {
726 self.recipients
727 .iter()
728 .position(|member| member == recipient)
729 }
730}
731
732impl<A, Route> Router<A, Route, LeastLoaded>
733where
734 A: Address,
735 Route: DeliveryRoute<Protocol: Protocol<Addr = A>> + Clone + PartialEq,
736{
737 #[must_use]
742 pub fn load_evidence(&self, recipient: &Route) -> Option<LoadEvidence> {
743 self.member_index(recipient)
744 .and_then(|index| self.strategy.loads.get(index).copied())
745 }
746}
747
748impl<A, Route, K> Router<A, Route, ConsistentHash<K>>
749where
750 A: Address,
751 Route: DeliveryRoute<Protocol: Protocol<Addr = A>> + Clone + PartialEq,
752 <Route::Protocol as Protocol>::Msg: RouteKey<K>,
753{
754 #[must_use]
757 pub fn member_token_evidence(&self, recipient: &Route) -> Option<MemberTokenEvidence> {
758 self.member_index(recipient)
759 .and_then(|index| self.strategy.membership.evidence.get(index).copied())
760 }
761}
762
763impl<A, Route, K> Router<A, Route, RendezvousHash<K>>
764where
765 A: Address,
766 Route: DeliveryRoute<Protocol: Protocol<Addr = A>> + Clone + PartialEq,
767 <Route::Protocol as Protocol>::Msg: RouteKey<K>,
768{
769 #[must_use]
772 pub fn member_token_evidence(&self, recipient: &Route) -> Option<MemberTokenEvidence> {
773 self.member_index(recipient)
774 .and_then(|index| self.strategy.membership.evidence.get(index).copied())
775 }
776}
777
778impl<A, Route, R> BehaviorBase for Router<A, Route, R>
779where
780 A: Address,
781 Route: DeliveryRoute<Protocol: Protocol<Addr = A>> + Clone + PartialEq,
782 R: RoutingStrategy<Route>,
783{
784 type Base = Self;
785
786 fn base(&self) -> &Self::Base {
787 self
788 }
789}
790
791impl<A, Route, R> behavior::Protocol for Router<A, Route, R>
792where
793 A: Address,
794 Route: DeliveryRoute<Protocol: Protocol<Addr = A>> + Clone + PartialEq,
795 R: RoutingStrategy<Route>,
796{
797 type Addr = A;
798 type Msg = RouterMessage<Route, R>;
799}
800
801impl<A, Route, R> Behavior for Router<A, Route, R>
802where
803 A: Address,
804 Route: DeliveryRoute<Protocol: Protocol<Addr = A>> + Clone + PartialEq,
805 R: RoutingStrategy<Route>,
806 Route::Sends: behavior::SendsFor<User<A, RouterMessage<Route, R>>>,
807{
808 type Protocol = Self;
809 type Event = User<A, RouterMessage<Route, R>>;
810 type Sends = Route::Sends;
811 type Ph = Never;
812 type Error = RouterError<<Route::Protocol as Protocol>::Msg, R::Observation, R::Error>;
813 type Birth = NoBirths;
814
815 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
816 match event.message {
817 RouterMessage::Add(recipient) => {
818 if !self.recipients.contains(&recipient) {
819 self.recipients.push(recipient.clone());
820 self.strategy.added(recipient);
821 }
822 Ok(Actions::cont())
823 }
824 RouterMessage::Remove(recipient) => {
825 if let Some(index) = self.recipients.iter().position(|item| item == &recipient) {
826 let removed = self.recipients.remove(index);
827 self.strategy.removed(index, removed, self.recipients.len());
828 }
829 Ok(Actions::cont())
830 }
831 RouterMessage::Route(message) => {
832 let mut strategy = self.strategy.clone();
833 let Some(index) = strategy.select(&self.recipients, &message) else {
834 return Err(RouterError::NoEligibleRecipients(message));
835 };
836 if index >= self.recipients.len() {
837 return Err(RouterError::InvalidSelection {
838 message,
839 index,
840 members: self.recipients.len(),
841 });
842 }
843 let sends = self.recipients[index].clone().deliver(message);
844 self.strategy = strategy;
845 Ok(Actions::send(sends))
846 }
847 RouterMessage::Observe(observation) => {
848 let mut strategy = self.strategy.clone();
849 strategy.observe(&self.recipients, observation).map_err(
850 |RoutingObservationRejection {
851 observation,
852 reason,
853 }| RouterError::Policy {
854 observation,
855 error: reason,
856 },
857 )?;
858 self.strategy = strategy;
859 Ok(Actions::cont())
860 }
861 }
862 }
863}
864
865#[cfg(test)]
866mod tests {
867 use super::*;
868 use crate::Activate as _;
869 use behavior::{Delivery, MailAddr, Recipient, Step};
870
871 struct Destination;
872
873 #[derive(Debug, Clone, PartialEq, Eq)]
874 struct KeyedMessage {
875 key: Key,
876 value: u8,
877 }
878
879 #[derive(Debug, Clone, PartialEq, Eq)]
880 struct Key(u64);
881
882 impl RouteKey<Key> for KeyedMessage {
883 fn route_key(&self) -> &Key {
884 &self.key
885 }
886 }
887
888 struct KeyedDestination;
889
890 impl behavior::Protocol for Destination {
891 type Addr = MailAddr;
892 type Msg = u8;
893 }
894
895 impl Behavior for Destination {
896 type Protocol = Self;
897 type Event = User<MailAddr, u8>;
898 type Sends = Vec<Never>;
899 type Ph = Never;
900 type Error = Never;
901 type Birth = NoBirths;
902
903 fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
904 Ok(Actions::cont())
905 }
906 }
907
908 impl behavior::Protocol for KeyedDestination {
909 type Addr = MailAddr;
910 type Msg = KeyedMessage;
911 }
912
913 impl Behavior for KeyedDestination {
914 type Protocol = Self;
915 type Event = User<MailAddr, KeyedMessage>;
916 type Sends = Vec<Never>;
917 type Ph = Never;
918 type Error = Never;
919 type Birth = NoBirths;
920
921 fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
922 Ok(Actions::cont())
923 }
924 }
925
926 #[test]
927 fn round_robin_repairs_cursor_after_removal() {
928 let one = Recipient::<Destination>::global(MailAddr(1));
929 let two = Recipient::<Destination>::global(MailAddr(2));
930 let three = Recipient::<Destination>::global(MailAddr(3));
931 let mut router = (Router::new(vec![one, two, three], RoundRobin::default()))
932 .initialize()
933 .unwrap()
934 .behavior;
935
936 let first = router
937 .receive(MailAddr(9), RouterMessage::Route(7))
938 .unwrap();
939 assert!(first.sends == vec![Delivery::new(one, 7)]);
940 assert!(matches!(first.become_, Step::Continue));
941
942 let removed = router
943 .receive(MailAddr(9), RouterMessage::Remove(one))
944 .unwrap();
945 assert!(removed.sends.is_empty());
946 assert!(removed.creates.is_empty());
947 assert_eq!(removed.become_, Step::Continue);
948 let second = router
949 .receive(MailAddr(9), RouterMessage::Route(8))
950 .unwrap();
951 assert!(second.sends == vec![Delivery::new(two, 8)]);
952 }
953
954 #[test]
955 fn empty_membership_returns_the_owned_payload() {
956 let mut router =
957 (Router::<MailAddr, Recipient<Destination>, _>::new(Vec::new(), RoundRobin::default()))
958 .initialize()
959 .unwrap()
960 .behavior;
961
962 let rejection = router.receive(MailAddr(9), RouterMessage::Route(11));
963 assert!(matches!(
964 rejection,
965 Err(RouterError::NoEligibleRecipients(11))
966 ));
967 assert!(router.recipients().is_empty());
968 }
969
970 #[test]
971 fn least_loaded_requires_typed_evidence_and_breaks_ties_by_membership_order() {
972 let one = Recipient::<Destination>::global(MailAddr(1));
973 let two = Recipient::<Destination>::global(MailAddr(2));
974 let mut router = (Router::new(vec![one, two], LeastLoaded::new()))
975 .initialize()
976 .unwrap()
977 .behavior;
978
979 let rejection = router.receive(MailAddr(9), RouterMessage::Route(1));
980 assert!(matches!(
981 rejection,
982 Err(RouterError::NoEligibleRecipients(1))
983 ));
984 for recipient in [one, two] {
985 let observed = router
986 .receive(
987 MailAddr(9),
988 RouterMessage::Observe(LoadObservation {
989 recipient,
990 version: LoadVersion(0),
991 load: Load(3),
992 }),
993 )
994 .unwrap();
995 assert!(observed.sends.is_empty());
996 assert!(observed.creates.is_empty());
997 assert_eq!(observed.become_, Step::Continue);
998 }
999 let tied = router
1000 .receive(MailAddr(9), RouterMessage::Route(2))
1001 .unwrap();
1002 assert!(tied.sends == vec![Delivery::new(one, 2)]);
1003
1004 let observed = router
1005 .receive(
1006 MailAddr(9),
1007 RouterMessage::Observe(LoadObservation {
1008 recipient: two,
1009 version: LoadVersion(1),
1010 load: Load(1),
1011 }),
1012 )
1013 .unwrap();
1014 assert!(observed.sends.is_empty());
1015 assert!(observed.creates.is_empty());
1016 assert_eq!(observed.become_, Step::Continue);
1017 let selected = router
1018 .receive(MailAddr(9), RouterMessage::Route(3))
1019 .unwrap();
1020 assert!(selected.sends == vec![Delivery::new(two, 3)]);
1021 }
1022
1023 #[test]
1024 fn least_loaded_rejects_stale_and_unknown_evidence_without_mutation() {
1025 let one = Recipient::<Destination>::global(MailAddr(1));
1026 let unknown = Recipient::<Destination>::global(MailAddr(8));
1027 let mut router = (Router::new(vec![one], LeastLoaded::new()))
1028 .initialize()
1029 .unwrap()
1030 .behavior;
1031 let observed = router
1032 .receive(
1033 MailAddr(9),
1034 RouterMessage::Observe(LoadObservation {
1035 recipient: one,
1036 version: LoadVersion(2),
1037 load: Load(4),
1038 }),
1039 )
1040 .unwrap();
1041 assert!(observed.sends.is_empty());
1042 assert!(observed.creates.is_empty());
1043 assert_eq!(observed.become_, Step::Continue);
1044
1045 let rejection = router.receive(
1046 MailAddr(9),
1047 RouterMessage::Observe(LoadObservation {
1048 recipient: one,
1049 version: LoadVersion(1),
1050 load: Load(0),
1051 }),
1052 );
1053 assert!(matches!(
1054 rejection,
1055 Err(RouterError::Policy {
1056 error: MemberEvidenceError::Stale,
1057 ..
1058 })
1059 ));
1060 let rejection = router.receive(
1061 MailAddr(9),
1062 RouterMessage::Observe(LoadObservation {
1063 recipient: one,
1064 version: LoadVersion(2),
1065 load: Load(5),
1066 }),
1067 );
1068 assert!(matches!(
1069 rejection,
1070 Err(RouterError::Policy {
1071 error: MemberEvidenceError::ConflictingVersion,
1072 ..
1073 })
1074 ));
1075 let rejection = router.receive(
1076 MailAddr(9),
1077 RouterMessage::Observe(LoadObservation {
1078 recipient: unknown,
1079 version: LoadVersion(0),
1080 load: Load(0),
1081 }),
1082 );
1083 assert!(matches!(
1084 rejection,
1085 Err(RouterError::Policy {
1086 error: MemberEvidenceError::UnknownRecipient,
1087 ..
1088 })
1089 ));
1090 assert_eq!(
1091 router.load_evidence(&one),
1092 Some(LoadEvidence::Observed {
1093 version: LoadVersion(2),
1094 load: Load(4)
1095 })
1096 );
1097 }
1098
1099 #[test]
1100 fn least_loaded_returns_unknown_observation_for_untracked_member() {
1101 let recipient = Recipient::<Destination>::global(MailAddr(1));
1102 let observation = LoadObservation {
1103 recipient,
1104 version: LoadVersion(4),
1105 load: Load(7),
1106 };
1107 let mut policy = LeastLoaded::new();
1108
1109 let rejection = policy.observe(&[recipient], observation);
1110 assert!(matches!(
1111 rejection,
1112 Err(RoutingObservationRejection {
1113 observation: LoadObservation {
1114 recipient: returned,
1115 version: LoadVersion(4),
1116 load: Load(7),
1117 },
1118 reason: MemberEvidenceError::UnknownRecipient,
1119 }) if returned == recipient
1120 ));
1121 }
1122
1123 fn identity_hash(key: &Key) -> u64 {
1124 key.0
1125 }
1126
1127 #[test]
1128 fn hash_policies_return_untracked_member_evidence() {
1129 let recipient = Recipient::<KeyedDestination>::global(MailAddr(1));
1130 let observation = MemberTokenObservation {
1131 recipient,
1132 version: MemberTokenVersion(3),
1133 token: MemberToken(8),
1134 };
1135 let mut ring = ConsistentHash::new(NonZeroU16::new(1).unwrap(), identity_hash);
1136 let mut rendezvous = RendezvousHash::new(identity_hash);
1137
1138 for rejection in [
1139 ring.observe(&[recipient], observation.clone()),
1140 rendezvous.observe(&[recipient], observation),
1141 ] {
1142 assert!(matches!(
1143 rejection,
1144 Err(RoutingObservationRejection {
1145 observation: MemberTokenObservation {
1146 recipient: returned,
1147 version: MemberTokenVersion(3),
1148 token: MemberToken(8),
1149 },
1150 reason: MemberEvidenceError::UnknownRecipient,
1151 }) if returned == recipient
1152 ));
1153 }
1154 }
1155
1156 #[test]
1157 fn consistent_hash_removal_moves_only_keys_owned_by_the_removed_member() {
1158 let members = [
1159 Recipient::<KeyedDestination>::global(MailAddr(1)),
1160 Recipient::<KeyedDestination>::global(MailAddr(2)),
1161 Recipient::<KeyedDestination>::global(MailAddr(3)),
1162 ];
1163 let mut router = (Router::new(
1164 members.to_vec(),
1165 ConsistentHash::new(NonZeroU16::new(8).unwrap(), identity_hash),
1166 ))
1167 .initialize()
1168 .unwrap()
1169 .behavior;
1170 for (index, recipient) in members.into_iter().enumerate() {
1171 let observed = router
1172 .receive(
1173 MailAddr(9),
1174 RouterMessage::Observe(MemberTokenObservation {
1175 recipient,
1176 version: MemberTokenVersion(0),
1177 token: MemberToken(u64::try_from(index + 1).unwrap()),
1178 }),
1179 )
1180 .unwrap();
1181 assert!(observed.sends.is_empty());
1182 assert!(observed.creates.is_empty());
1183 assert_eq!(observed.become_, Step::Continue);
1184 }
1185 let before = (0..128_u64)
1186 .map(|key| {
1187 router
1188 .receive(
1189 MailAddr(9),
1190 RouterMessage::Route(KeyedMessage {
1191 key: Key(key),
1192 value: 1,
1193 }),
1194 )
1195 .unwrap()
1196 .sends[0]
1197 .to
1198 })
1199 .collect::<Vec<_>>();
1200 let removed = router
1201 .receive(MailAddr(9), RouterMessage::Remove(members[1]))
1202 .unwrap();
1203 assert!(removed.sends.is_empty());
1204 assert!(removed.creates.is_empty());
1205 assert_eq!(removed.become_, Step::Continue);
1206 for (key, previous) in before.into_iter().enumerate() {
1207 let current = router
1208 .receive(
1209 MailAddr(9),
1210 RouterMessage::Route(KeyedMessage {
1211 key: Key(u64::try_from(key).unwrap()),
1212 value: 1,
1213 }),
1214 )
1215 .unwrap()
1216 .sends[0]
1217 .to;
1218 if previous != members[1] {
1219 assert!(current == previous);
1220 }
1221 }
1222 }
1223
1224 #[test]
1225 fn rendezvous_hash_is_deterministic_and_rejects_conflicting_tokens() {
1226 let one = Recipient::<KeyedDestination>::global(MailAddr(1));
1227 let two = Recipient::<KeyedDestination>::global(MailAddr(2));
1228 let mut router = (Router::new(vec![one, two], RendezvousHash::new(identity_hash)))
1229 .initialize()
1230 .unwrap()
1231 .behavior;
1232 for (recipient, token) in [(one, 11), (two, 22)] {
1233 let observed = router
1234 .receive(
1235 MailAddr(9),
1236 RouterMessage::Observe(MemberTokenObservation {
1237 recipient,
1238 version: MemberTokenVersion(0),
1239 token: MemberToken(token),
1240 }),
1241 )
1242 .unwrap();
1243 assert!(observed.sends.is_empty());
1244 assert!(observed.creates.is_empty());
1245 assert_eq!(observed.become_, Step::Continue);
1246 }
1247 let first = router
1248 .receive(
1249 MailAddr(9),
1250 RouterMessage::Route(KeyedMessage {
1251 key: Key(7),
1252 value: 1,
1253 }),
1254 )
1255 .unwrap()
1256 .sends[0]
1257 .to;
1258 let again = router
1259 .receive(
1260 MailAddr(9),
1261 RouterMessage::Route(KeyedMessage {
1262 key: Key(7),
1263 value: 2,
1264 }),
1265 )
1266 .unwrap()
1267 .sends[0]
1268 .to;
1269 assert!(first == again);
1270 let conflicting = MemberTokenObservation {
1271 recipient: one,
1272 version: MemberTokenVersion(0),
1273 token: MemberToken(99),
1274 };
1275 let rejection = router.receive(MailAddr(9), RouterMessage::Observe(conflicting));
1276 assert!(matches!(
1277 rejection,
1278 Err(RouterError::Policy {
1279 observation: MemberTokenObservation {
1280 recipient: returned,
1281 version: MemberTokenVersion(0),
1282 token: MemberToken(99),
1283 },
1284 error: MemberEvidenceError::ConflictingVersion,
1285 }) if returned == one
1286 ));
1287 assert_eq!(
1288 router.member_token_evidence(&one),
1289 Some(MemberTokenEvidence::Observed {
1290 version: MemberTokenVersion(0),
1291 token: MemberToken(11)
1292 })
1293 );
1294 }
1295
1296 #[test]
1297 fn hash_token_versions_return_stale_evidence_and_accept_newer_observations() {
1298 let recipient = Recipient::<KeyedDestination>::global(MailAddr(1));
1299 let mut router = (Router::new(vec![recipient], RendezvousHash::new(identity_hash)))
1300 .initialize()
1301 .unwrap()
1302 .behavior;
1303 let observed = router
1304 .receive(
1305 MailAddr(9),
1306 RouterMessage::Observe(MemberTokenObservation {
1307 recipient,
1308 version: MemberTokenVersion(2),
1309 token: MemberToken(11),
1310 }),
1311 )
1312 .unwrap();
1313 assert!(observed.sends.is_empty());
1314 assert!(observed.creates.is_empty());
1315 assert_eq!(observed.become_, Step::Continue);
1316
1317 let stale = MemberTokenObservation {
1318 recipient,
1319 version: MemberTokenVersion(1),
1320 token: MemberToken(99),
1321 };
1322 let rejection = router.receive(MailAddr(9), RouterMessage::Observe(stale));
1323 assert!(matches!(
1324 rejection,
1325 Err(RouterError::Policy {
1326 observation: MemberTokenObservation {
1327 recipient: returned,
1328 version: MemberTokenVersion(1),
1329 token: MemberToken(99),
1330 },
1331 error: MemberEvidenceError::Stale,
1332 }) if returned == recipient
1333 ));
1334 assert_eq!(
1335 router.member_token_evidence(&recipient),
1336 Some(MemberTokenEvidence::Observed {
1337 version: MemberTokenVersion(2),
1338 token: MemberToken(11),
1339 })
1340 );
1341
1342 for (version, token) in [(2, 11), (3, 22)] {
1343 let accepted = router
1344 .receive(
1345 MailAddr(9),
1346 RouterMessage::Observe(MemberTokenObservation {
1347 recipient,
1348 version: MemberTokenVersion(version),
1349 token: MemberToken(token),
1350 }),
1351 )
1352 .unwrap();
1353 assert!(accepted.sends.is_empty());
1354 assert!(accepted.creates.is_empty());
1355 assert_eq!(accepted.become_, Step::Continue);
1356 }
1357 assert_eq!(
1358 router.member_token_evidence(&recipient),
1359 Some(MemberTokenEvidence::Observed {
1360 version: MemberTokenVersion(3),
1361 token: MemberToken(22),
1362 })
1363 );
1364 }
1365
1366 enum CapacityReading {
1367 Available(Box<u64>),
1368 Unavailable(Box<u64>),
1369 }
1370
1371 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
1372 enum CapacityRefusal {
1373 NoReading,
1374 }
1375
1376 #[derive(Clone, Default)]
1377 struct CapacityPolicy {
1378 current: Option<u64>,
1379 }
1380
1381 impl RoutingStrategy<Recipient<Destination>> for CapacityPolicy {
1382 type Observation = CapacityReading;
1383 type Error = CapacityRefusal;
1384
1385 fn select(&mut self, members: &[Recipient<Destination>], _: &u8) -> Option<usize> {
1386 if self.current.is_some() && !members.is_empty() {
1387 Some(0)
1388 } else {
1389 None
1390 }
1391 }
1392
1393 fn observe(
1394 &mut self,
1395 _: &[Recipient<Destination>],
1396 observation: Self::Observation,
1397 ) -> Result<(), RoutingObservationRejection<Self::Observation, Self::Error>> {
1398 match observation {
1399 CapacityReading::Available(amount) => {
1400 self.current = Some(*amount);
1401 Ok(())
1402 }
1403 observation @ CapacityReading::Unavailable(_) => {
1404 self.current = Some(999);
1405 Err(RoutingObservationRejection {
1406 observation,
1407 reason: CapacityRefusal::NoReading,
1408 })
1409 }
1410 }
1411 }
1412 }
1413
1414 #[test]
1415 fn rejected_noncloning_observation_returns_original_allocation_and_policy_state() {
1416 let recipient = Recipient::<Destination>::global(MailAddr(1));
1417 let mut router = (Router::new(vec![recipient], CapacityPolicy::default()))
1418 .initialize()
1419 .unwrap()
1420 .behavior;
1421 let reading = Box::new(13);
1422 let original = (&*reading) as *const u64;
1423 let rejection = router.receive(
1424 MailAddr(9),
1425 RouterMessage::Observe(CapacityReading::Unavailable(reading)),
1426 );
1427 match rejection {
1428 Err(RouterError::Policy {
1429 observation: CapacityReading::Unavailable(returned),
1430 error: CapacityRefusal::NoReading,
1431 }) => {
1432 assert_eq!((&*returned) as *const u64, original);
1433 assert_eq!(*returned, 13);
1434 }
1435 _ => panic!("unavailable reading must return its original allocation"),
1436 }
1437 assert_eq!(router.strategy().current, None);
1438
1439 let observed = router
1440 .receive(
1441 MailAddr(9),
1442 RouterMessage::Observe(CapacityReading::Available(Box::new(3))),
1443 )
1444 .unwrap();
1445 assert!(observed.sends.is_empty());
1446 assert!(observed.creates.is_empty());
1447 assert_eq!(observed.become_, Step::Continue);
1448 assert_eq!(router.strategy().current, Some(3));
1449
1450 let routed = router
1451 .receive(MailAddr(9), RouterMessage::Route(42))
1452 .unwrap();
1453 assert!(routed.sends == vec![Delivery::new(recipient, 42)]);
1454 assert!(routed.creates.is_empty());
1455 assert_eq!(routed.become_, Step::Continue);
1456 }
1457
1458 #[derive(Clone, Default)]
1459 struct RejectAfterMutation {
1460 selections: usize,
1461 observations: usize,
1462 }
1463
1464 impl RoutingStrategy<Recipient<Destination>> for RejectAfterMutation {
1465 type Observation = u8;
1466 type Error = u8;
1467
1468 fn select(&mut self, members: &[Recipient<Destination>], _: &u8) -> Option<usize> {
1469 self.selections += 1;
1470 Some(members.len())
1471 }
1472
1473 fn observe(
1474 &mut self,
1475 _: &[Recipient<Destination>],
1476 observation: Self::Observation,
1477 ) -> Result<(), RoutingObservationRejection<Self::Observation, Self::Error>> {
1478 self.observations += 1;
1479 Err(RoutingObservationRejection {
1480 observation,
1481 reason: 99,
1482 })
1483 }
1484 }
1485
1486 #[test]
1487 fn rejected_policy_turns_preserve_the_command_and_policy_snapshot() {
1488 let member = Recipient::<Destination>::global(MailAddr(1));
1489 let mut router = Router::new(vec![member], RejectAfterMutation::default())
1490 .initialize()
1491 .unwrap()
1492 .behavior;
1493
1494 let rejection = router.receive(MailAddr(9), RouterMessage::Route(42));
1495 assert!(matches!(
1496 rejection,
1497 Err(RouterError::InvalidSelection {
1498 message: 42,
1499 index: 1,
1500 members: 1,
1501 })
1502 ));
1503 assert_eq!(router.strategy().selections, 0);
1504
1505 let rejection = router.receive(MailAddr(9), RouterMessage::Observe(7));
1506 assert!(matches!(
1507 rejection,
1508 Err(RouterError::Policy {
1509 observation: 7,
1510 error: 99,
1511 })
1512 ));
1513 assert_eq!(router.strategy().observations, 0);
1514 }
1515}