1use std::collections::{BTreeMap, BTreeSet, VecDeque};
11use std::sync::{Arc, Mutex};
12
13use super::MemberId;
14
15#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
17pub enum LivenessStatus {
18 Alive,
19 Suspect,
20 Unreachable,
21}
22
23#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
25pub enum LoadBucket {
26 Low,
27 Medium,
28 High,
29 Saturated,
30}
31
32#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
34pub struct LivenessObservation {
35 pub observer: MemberId,
36 pub observed: MemberId,
37 pub incarnation: u64,
38 pub status: LivenessStatus,
39}
40
41#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
43pub struct MemberHealthInput {
44 pub member: MemberId,
45 pub error_count: u32,
46 pub replication_lag_records: u64,
47 pub read_only: bool,
48 pub self_fenced: bool,
49}
50
51#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
53pub struct LoadMetricSample {
54 pub member: MemberId,
55 pub disk_pressure: LoadBucket,
56 pub cpu_pressure: LoadBucket,
57 pub range_hotness: LoadBucket,
58 pub write_throughput_bucket: u16,
59 pub read_throughput_bucket: u16,
60}
61
62#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
64pub struct CatalogVersionHint {
65 pub member: MemberId,
66 pub ownership_catalog_version: u64,
67 pub topology_generation: u64,
68 pub placement_generation: u64,
69}
70
71#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
73pub struct TopologyHint {
74 pub member: MemberId,
75 pub endpoint: String,
76 pub region: String,
77 pub failure_domain: String,
78}
79
80#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
86pub enum SignalPlaneMessage {
87 LivenessObservation(LivenessObservation),
88 MemberHealthInput(MemberHealthInput),
89 LoadMetricSample(LoadMetricSample),
90 CatalogVersionHint(CatalogVersionHint),
91 TopologyHint(TopologyHint),
92}
93
94#[derive(Debug, Clone, PartialEq, Eq)]
96pub struct ReceivedSignal {
97 pub from: MemberId,
98 pub to: MemberId,
99 pub message: SignalPlaneMessage,
100}
101
102pub trait SignalPlane {
108 fn set_members(&mut self, members: Vec<MemberId>);
110
111 fn publish(&mut self, from: MemberId, message: SignalPlaneMessage);
113
114 fn drain_received(&mut self, member: &MemberId) -> Vec<ReceivedSignal>;
116}
117
118#[derive(Debug, Clone, Copy, PartialEq, Eq)]
120pub struct SignalPlaneLimits {
121 pub fanout: usize,
123 pub max_payload_messages: usize,
125}
126
127impl Default for SignalPlaneLimits {
128 fn default() -> Self {
129 Self {
130 fanout: 3,
131 max_payload_messages: 16,
132 }
133 }
134}
135
136#[derive(Debug, Clone, Default, PartialEq, Eq)]
138pub struct SignalPlaneMetrics {
139 pub sent_frames_total: u64,
140 pub sent_messages_total: u64,
141 pub received_frames_total: u64,
142 pub received_messages_total: u64,
143 pub rejected_frames_total: u64,
144 pub dropped_frames_total: u64,
145 pub payload_cap_hits_total: u64,
146 pub fanout_cap: usize,
147 pub payload_cap: usize,
148}
149
150#[derive(Debug, Clone, PartialEq, Eq)]
152pub struct SignalFrame {
153 pub authenticated_peer: MemberId,
154 pub from: MemberId,
155 pub to: MemberId,
156 pub messages: Vec<SignalPlaneMessage>,
157}
158
159#[derive(Debug, Clone, PartialEq, Eq)]
161pub enum SignalTransportError {
162 AuthenticatedPeerMismatch {
163 authenticated_peer: MemberId,
164 declared_sender: MemberId,
165 },
166 SenderNotAdmitted(MemberId),
167 RecipientNotAdmitted(MemberId),
168 PayloadTooLarge {
169 actual: usize,
170 max: usize,
171 },
172 EndpointUnavailable(MemberId),
173}
174
175#[derive(Debug)]
177pub struct IntraClusterSignalBus {
178 admitted: BTreeSet<MemberId>,
179 endpoints: BTreeMap<MemberId, VecDeque<SignalFrame>>,
180 limits: SignalPlaneLimits,
181 metrics: SignalPlaneMetrics,
182}
183
184impl Default for IntraClusterSignalBus {
185 fn default() -> Self {
186 Self::new(SignalPlaneLimits::default())
187 }
188}
189
190impl IntraClusterSignalBus {
191 pub fn new(limits: SignalPlaneLimits) -> Self {
192 let metrics = SignalPlaneMetrics {
193 fanout_cap: limits.fanout,
194 payload_cap: limits.max_payload_messages,
195 ..Default::default()
196 };
197 Self {
198 admitted: BTreeSet::new(),
199 endpoints: BTreeMap::new(),
200 limits,
201 metrics,
202 }
203 }
204
205 pub fn set_limits(&mut self, limits: SignalPlaneLimits) {
206 self.limits = limits;
207 self.metrics.fanout_cap = limits.fanout;
208 self.metrics.payload_cap = limits.max_payload_messages;
209 }
210
211 pub fn set_admitted_members(&mut self, members: Vec<MemberId>) {
212 self.admitted = members.into_iter().collect();
213 self.endpoints
214 .retain(|member, _| self.admitted.contains(member));
215 }
216
217 pub fn register_endpoint(&mut self, member: MemberId) {
218 if self.admitted.is_empty() || self.admitted.contains(&member) {
219 self.endpoints.entry(member).or_default();
220 }
221 }
222
223 pub fn unregister_endpoint(&mut self, member: &MemberId) {
224 self.endpoints.remove(member);
225 }
226
227 pub fn send(&mut self, frame: SignalFrame) -> Result<(), SignalTransportError> {
228 if frame.authenticated_peer != frame.from {
229 self.metrics.rejected_frames_total += 1;
230 return Err(SignalTransportError::AuthenticatedPeerMismatch {
231 authenticated_peer: frame.authenticated_peer,
232 declared_sender: frame.from,
233 });
234 }
235 if !self.admitted.contains(&frame.from) {
236 self.metrics.rejected_frames_total += 1;
237 return Err(SignalTransportError::SenderNotAdmitted(frame.from));
238 }
239 if !self.admitted.contains(&frame.to) {
240 self.metrics.rejected_frames_total += 1;
241 return Err(SignalTransportError::RecipientNotAdmitted(frame.to));
242 }
243 if frame.messages.len() > self.limits.max_payload_messages {
244 self.metrics.rejected_frames_total += 1;
245 self.metrics.payload_cap_hits_total += 1;
246 return Err(SignalTransportError::PayloadTooLarge {
247 actual: frame.messages.len(),
248 max: self.limits.max_payload_messages,
249 });
250 }
251
252 let message_count = frame.messages.len() as u64;
253 let Some(endpoint) = self.endpoints.get_mut(&frame.to) else {
254 self.metrics.dropped_frames_total += 1;
255 return Err(SignalTransportError::EndpointUnavailable(frame.to));
256 };
257 endpoint.push_back(frame);
258 self.metrics.sent_frames_total += 1;
259 self.metrics.sent_messages_total += message_count;
260 Ok(())
261 }
262
263 pub fn drain(&mut self, member: MemberId) -> Vec<SignalFrame> {
264 let Some(endpoint) = self.endpoints.get_mut(&member) else {
265 return Vec::new();
266 };
267 endpoint.drain(..).collect()
268 }
269
270 pub fn metrics(&self) -> SignalPlaneMetrics {
271 self.metrics.clone()
272 }
273}
274
275#[derive(Debug, Clone, Default)]
277pub struct SharedSignalTransport {
278 bus: Arc<Mutex<IntraClusterSignalBus>>,
279}
280
281impl SharedSignalTransport {
282 pub fn set_limits(&self, limits: SignalPlaneLimits) {
283 self.bus
284 .lock()
285 .expect("signal transport mutex poisoned")
286 .set_limits(limits);
287 }
288
289 pub fn set_admitted_members(&self, members: Vec<MemberId>) {
290 self.bus
291 .lock()
292 .expect("signal transport mutex poisoned")
293 .set_admitted_members(members);
294 }
295
296 pub fn register_endpoint(&self, member: MemberId) {
297 self.bus
298 .lock()
299 .expect("signal transport mutex poisoned")
300 .register_endpoint(member);
301 }
302
303 pub fn unregister_endpoint(&self, member: &MemberId) {
304 self.bus
305 .lock()
306 .expect("signal transport mutex poisoned")
307 .unregister_endpoint(member);
308 }
309
310 pub fn send(&self, frame: SignalFrame) -> Result<(), SignalTransportError> {
311 self.bus
312 .lock()
313 .expect("signal transport mutex poisoned")
314 .send(frame)
315 }
316
317 pub fn drain(&self, member: MemberId) -> Vec<SignalFrame> {
318 self.bus
319 .lock()
320 .expect("signal transport mutex poisoned")
321 .drain(member)
322 }
323
324 pub fn metrics(&self) -> SignalPlaneMetrics {
325 self.bus
326 .lock()
327 .expect("signal transport mutex poisoned")
328 .metrics()
329 }
330}
331
332#[derive(Debug)]
334pub struct TransportSignalPlane {
335 local_member: MemberId,
336 transport: SharedSignalTransport,
337 limits: SignalPlaneLimits,
338 members: BTreeSet<MemberId>,
339 known: BTreeSet<SignalPlaneMessage>,
340 inbox: Vec<ReceivedSignal>,
341 round: u64,
342 metrics: SignalPlaneMetrics,
343}
344
345impl TransportSignalPlane {
346 pub fn new(
347 local_member: MemberId,
348 transport: SharedSignalTransport,
349 limits: SignalPlaneLimits,
350 ) -> Self {
351 transport.set_limits(limits);
352 transport.register_endpoint(local_member.clone());
353 let metrics = SignalPlaneMetrics {
354 fanout_cap: limits.fanout,
355 payload_cap: limits.max_payload_messages,
356 ..Default::default()
357 };
358 Self {
359 local_member,
360 transport,
361 limits,
362 members: BTreeSet::new(),
363 known: BTreeSet::new(),
364 inbox: Vec::new(),
365 round: 0,
366 metrics,
367 }
368 }
369
370 pub fn advance_round(&mut self) {
371 self.receive_frames();
372 self.send_round();
373 self.round += 1;
374 }
375
376 pub fn knows(&self, message: &SignalPlaneMessage) -> bool {
377 self.known.contains(message)
378 }
379
380 pub fn metrics(&self) -> SignalPlaneMetrics {
381 self.metrics.clone()
382 }
383
384 fn receive_frames(&mut self) {
385 let frames = self.transport.drain(self.local_member.clone());
386 for frame in frames {
387 if !self.members.contains(&frame.from) || !self.members.contains(&frame.to) {
388 self.metrics.rejected_frames_total += 1;
389 continue;
390 }
391 self.metrics.received_frames_total += 1;
392 self.metrics.received_messages_total += frame.messages.len() as u64;
393 for message in frame.messages {
394 self.known.insert(message.clone());
395 self.inbox.push(ReceivedSignal {
396 from: frame.from.clone(),
397 to: self.local_member.clone(),
398 message,
399 });
400 }
401 }
402 }
403
404 fn send_round(&mut self) {
405 if self.known.is_empty() {
406 return;
407 }
408
409 let payload = self
410 .known
411 .iter()
412 .take(self.limits.max_payload_messages)
413 .cloned()
414 .collect::<Vec<_>>();
415 if self.known.len() > payload.len() {
416 self.metrics.payload_cap_hits_total += 1;
417 }
418
419 for peer in self.sample_peers() {
420 let message_count = payload.len() as u64;
421 let frame = SignalFrame {
422 authenticated_peer: self.local_member.clone(),
423 from: self.local_member.clone(),
424 to: peer,
425 messages: payload.clone(),
426 };
427 match self.transport.send(frame) {
428 Ok(()) => {
429 self.metrics.sent_frames_total += 1;
430 self.metrics.sent_messages_total += message_count;
431 }
432 Err(SignalTransportError::EndpointUnavailable(_)) => {
433 self.metrics.dropped_frames_total += 1;
434 }
435 Err(_) => {
436 self.metrics.rejected_frames_total += 1;
437 }
438 }
439 }
440 }
441
442 fn sample_peers(&self) -> Vec<MemberId> {
443 if self.limits.fanout == 0 {
444 return Vec::new();
445 }
446
447 let peers = self
448 .members
449 .iter()
450 .filter(|member| *member != &self.local_member)
451 .cloned()
452 .collect::<Vec<_>>();
453 if peers.len() <= self.limits.fanout {
454 return peers;
455 }
456
457 let start = (self.round + stable_member_hash(&self.local_member)) as usize % peers.len();
458 (0..self.limits.fanout)
459 .map(|offset| peers[(start + offset) % peers.len()].clone())
460 .collect()
461 }
462}
463
464impl SignalPlane for TransportSignalPlane {
465 fn set_members(&mut self, members: Vec<MemberId>) {
466 self.members = members.iter().cloned().collect();
467 self.transport.set_admitted_members(members);
468 if self.members.contains(&self.local_member) {
469 self.transport.register_endpoint(self.local_member.clone());
470 }
471 }
472
473 fn publish(&mut self, from: MemberId, message: SignalPlaneMessage) {
474 if from == self.local_member && self.members.contains(&from) {
475 self.known.insert(message);
476 } else {
477 self.metrics.rejected_frames_total += 1;
478 }
479 }
480
481 fn drain_received(&mut self, member: &MemberId) -> Vec<ReceivedSignal> {
482 if member == &self.local_member {
483 self.receive_frames();
484 std::mem::take(&mut self.inbox)
485 } else {
486 Vec::new()
487 }
488 }
489}
490
491fn stable_member_hash(member: &MemberId) -> u64 {
492 member.bytes().fold(0_u64, |hash, byte| {
493 hash.wrapping_mul(31).wrapping_add(u64::from(byte))
494 })
495}
496
497#[derive(Debug, Clone, Default)]
499pub struct SignalPlaneSchedule {
500 dropped: BTreeSet<ScheduledLink>,
501 delays: BTreeMap<ScheduledLink, u64>,
502 partitions: Vec<Partition>,
503}
504
505impl SignalPlaneSchedule {
506 pub fn drop_delivery(&mut self, round: u64, from: MemberId, to: MemberId) {
508 self.dropped.insert(ScheduledLink { round, from, to });
509 }
510
511 pub fn delay_delivery(&mut self, round: u64, from: MemberId, to: MemberId, extra_rounds: u64) {
513 self.delays
514 .insert(ScheduledLink { round, from, to }, extra_rounds);
515 }
516
517 pub fn partition_between(
519 &mut self,
520 start_round: u64,
521 end_round: u64,
522 a: MemberId,
523 b: MemberId,
524 ) {
525 self.partitions.push(Partition {
526 start_round,
527 end_round,
528 a,
529 b,
530 });
531 }
532
533 fn is_dropped(&self, round: u64, from: &MemberId, to: &MemberId) -> bool {
534 self.dropped.contains(&ScheduledLink {
535 round,
536 from: from.clone(),
537 to: to.clone(),
538 })
539 }
540
541 fn delay(&self, round: u64, from: &MemberId, to: &MemberId) -> u64 {
542 self.delays
543 .get(&ScheduledLink {
544 round,
545 from: from.clone(),
546 to: to.clone(),
547 })
548 .copied()
549 .unwrap_or(0)
550 }
551
552 fn is_partitioned(&self, round: u64, from: &MemberId, to: &MemberId) -> bool {
553 self.partitions
554 .iter()
555 .any(|partition| partition.contains(round) && partition.matches_members(from, to))
556 }
557}
558
559#[derive(Debug, Default)]
561pub struct SimulatedSignalPlane {
562 members: BTreeMap<MemberId, SimulatedPeer>,
563 pending: Vec<PendingSignal>,
564 schedule: SignalPlaneSchedule,
565 round: u64,
566}
567
568impl SimulatedSignalPlane {
569 pub fn new(members: Vec<MemberId>) -> Self {
570 let mut plane = Self::default();
571 plane.set_members(members);
572 plane
573 }
574
575 pub fn schedule_mut(&mut self) -> &mut SignalPlaneSchedule {
576 &mut self.schedule
577 }
578
579 pub fn round(&self) -> u64 {
580 self.round
581 }
582
583 pub fn advance_round(&mut self) {
585 self.queue_known_signals();
586 self.deliver_due_signals();
587 self.round += 1;
588 }
589
590 pub fn members_with_signal(&self, message: &SignalPlaneMessage) -> Vec<MemberId> {
592 self.members
593 .iter()
594 .filter(|(_, peer)| peer.known.contains(message))
595 .map(|(member, _)| member.clone())
596 .collect()
597 }
598
599 fn queue_known_signals(&mut self) {
600 let members = self.members.keys().cloned().collect::<Vec<_>>();
601 let mut queued = Vec::new();
602
603 for from in &members {
604 let Some(peer) = self.members.get(from) else {
605 continue;
606 };
607 for to in &members {
608 if from == to {
609 continue;
610 }
611 if self.schedule.is_partitioned(self.round, from, to)
612 || self.schedule.is_dropped(self.round, from, to)
613 {
614 continue;
615 }
616 let deliver_at = self.round + self.schedule.delay(self.round, from, to);
617 for message in &peer.known {
618 queued.push(PendingSignal {
619 deliver_at,
620 signal: ReceivedSignal {
621 from: from.clone(),
622 to: to.clone(),
623 message: message.clone(),
624 },
625 });
626 }
627 }
628 }
629
630 self.pending.extend(queued);
631 }
632
633 fn deliver_due_signals(&mut self) {
634 let mut pending = Vec::new();
635 for signal in self.pending.drain(..) {
636 if signal.deliver_at <= self.round {
637 if let Some(peer) = self.members.get_mut(&signal.signal.to) {
638 peer.known.insert(signal.signal.message.clone());
639 peer.inbox.push(signal.signal);
640 }
641 } else {
642 pending.push(signal);
643 }
644 }
645 self.pending = pending;
646 }
647}
648
649impl SignalPlane for SimulatedSignalPlane {
650 fn set_members(&mut self, members: Vec<MemberId>) {
651 let admitted = members.into_iter().collect::<BTreeSet<_>>();
652 self.members.retain(|member, _| admitted.contains(member));
653 for member in admitted {
654 self.members.entry(member).or_default();
655 }
656 }
657
658 fn publish(&mut self, from: MemberId, message: SignalPlaneMessage) {
659 if let Some(peer) = self.members.get_mut(&from) {
660 peer.known.insert(message);
661 }
662 }
663
664 fn drain_received(&mut self, member: &MemberId) -> Vec<ReceivedSignal> {
665 self.members
666 .get_mut(member)
667 .map(|peer| std::mem::take(&mut peer.inbox))
668 .unwrap_or_default()
669 }
670}
671
672#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
673struct ScheduledLink {
674 round: u64,
675 from: MemberId,
676 to: MemberId,
677}
678
679#[derive(Debug, Clone)]
680struct Partition {
681 start_round: u64,
682 end_round: u64,
683 a: MemberId,
684 b: MemberId,
685}
686
687impl Partition {
688 fn contains(&self, round: u64) -> bool {
689 self.start_round <= round && round <= self.end_round
690 }
691
692 fn matches_members(&self, from: &MemberId, to: &MemberId) -> bool {
693 (&self.a == from && &self.b == to) || (&self.a == to && &self.b == from)
694 }
695}
696
697#[derive(Debug, Default)]
698struct SimulatedPeer {
699 known: BTreeSet<SignalPlaneMessage>,
700 inbox: Vec<ReceivedSignal>,
701}
702
703#[derive(Debug, Clone)]
704struct PendingSignal {
705 deliver_at: u64,
706 signal: ReceivedSignal,
707}
708
709#[cfg(test)]
710mod tests {
711 use super::*;
712
713 fn member(id: &str) -> MemberId {
714 id.to_string()
715 }
716
717 fn liveness(observer: &str, observed: &str) -> SignalPlaneMessage {
718 SignalPlaneMessage::LivenessObservation(LivenessObservation {
719 observer: member(observer),
720 observed: member(observed),
721 incarnation: 1,
722 status: LivenessStatus::Alive,
723 })
724 }
725
726 #[test]
727 fn message_set_is_closed_over_signal_plane_only() {
728 let messages = [
729 SignalPlaneMessage::LivenessObservation(LivenessObservation {
730 observer: member("a"),
731 observed: member("b"),
732 incarnation: 1,
733 status: LivenessStatus::Alive,
734 }),
735 SignalPlaneMessage::MemberHealthInput(MemberHealthInput {
736 member: member("a"),
737 error_count: 0,
738 replication_lag_records: 0,
739 read_only: false,
740 self_fenced: false,
741 }),
742 SignalPlaneMessage::LoadMetricSample(LoadMetricSample {
743 member: member("a"),
744 disk_pressure: LoadBucket::Low,
745 cpu_pressure: LoadBucket::Medium,
746 range_hotness: LoadBucket::High,
747 write_throughput_bucket: 2,
748 read_throughput_bucket: 3,
749 }),
750 SignalPlaneMessage::CatalogVersionHint(CatalogVersionHint {
751 member: member("a"),
752 ownership_catalog_version: 7,
753 topology_generation: 8,
754 placement_generation: 9,
755 }),
756 SignalPlaneMessage::TopologyHint(TopologyHint {
757 member: member("a"),
758 endpoint: "redb://a".to_string(),
759 region: "local".to_string(),
760 failure_domain: "rack-1".to_string(),
761 }),
762 ];
763
764 for message in messages {
765 match message {
766 SignalPlaneMessage::LivenessObservation(_)
767 | SignalPlaneMessage::MemberHealthInput(_)
768 | SignalPlaneMessage::LoadMetricSample(_)
769 | SignalPlaneMessage::CatalogVersionHint(_)
770 | SignalPlaneMessage::TopologyHint(_) => {}
771 }
772 }
773 }
774
775 struct FakeSignalPlane {
776 members: Vec<MemberId>,
777 inbox: Vec<ReceivedSignal>,
778 }
779
780 impl SignalPlane for FakeSignalPlane {
781 fn set_members(&mut self, members: Vec<MemberId>) {
782 self.members = members;
783 }
784
785 fn publish(&mut self, from: MemberId, message: SignalPlaneMessage) {
786 for member in &self.members {
787 if member != &from {
788 self.inbox.push(ReceivedSignal {
789 from: from.clone(),
790 to: member.clone(),
791 message: message.clone(),
792 });
793 }
794 }
795 }
796
797 fn drain_received(&mut self, member: &MemberId) -> Vec<ReceivedSignal> {
798 let mut drained = Vec::new();
799 self.inbox.retain(|signal| {
800 if &signal.to == member {
801 drained.push(signal.clone());
802 false
803 } else {
804 true
805 }
806 });
807 drained
808 }
809 }
810
811 #[test]
812 fn fake_signal_plane_exercises_the_narrow_trait() {
813 let mut plane = FakeSignalPlane {
814 members: Vec::new(),
815 inbox: Vec::new(),
816 };
817 plane.set_members(vec![member("a"), member("b")]);
818
819 plane.publish(member("a"), liveness("a", "b"));
820
821 assert!(plane.drain_received(&member("a")).is_empty());
822 assert_eq!(plane.drain_received(&member("b")).len(), 1);
823 }
824
825 #[test]
826 fn simulated_peers_converge_with_loss_delay_and_temporary_partition() {
827 let mut plane =
828 SimulatedSignalPlane::new(vec![member("a"), member("b"), member("c"), member("d")]);
829 plane
830 .schedule_mut()
831 .drop_delivery(0, member("a"), member("c"));
832 plane
833 .schedule_mut()
834 .delay_delivery(1, member("b"), member("d"), 1);
835 plane
836 .schedule_mut()
837 .partition_between(0, 1, member("a"), member("d"));
838
839 let observation = liveness("a", "b");
840 plane.publish(member("a"), observation.clone());
841
842 let bound = 6;
843 for _ in 0..bound {
844 if plane.members_with_signal(&observation).len() == 4 {
845 break;
846 }
847 plane.advance_round();
848 }
849
850 assert_eq!(plane.members_with_signal(&observation).len(), 4);
851 assert!(
852 plane.round() <= bound,
853 "liveness observation did not converge within {bound} rounds"
854 );
855 }
856
857 #[test]
858 fn simulated_partition_limits_convergence_to_connected_members() {
859 let mut plane =
860 SimulatedSignalPlane::new(vec![member("a"), member("b"), member("c"), member("d")]);
861 plane
862 .schedule_mut()
863 .partition_between(0, 10, member("c"), member("d"));
864 plane
865 .schedule_mut()
866 .partition_between(0, 10, member("a"), member("d"));
867 plane
868 .schedule_mut()
869 .partition_between(0, 10, member("b"), member("d"));
870
871 let observation = liveness("a", "b");
872 plane.publish(member("a"), observation.clone());
873
874 for _ in 0..4 {
875 plane.advance_round();
876 }
877
878 assert_eq!(
879 plane.members_with_signal(&observation),
880 vec![member("a"), member("b"), member("c")]
881 );
882 }
883
884 #[test]
885 fn secured_transport_rejects_signal_from_unadmitted_peer_before_delivery() {
886 let mut bus = IntraClusterSignalBus::new(SignalPlaneLimits::default());
887 bus.set_admitted_members(vec![member("a"), member("b")]);
888 bus.register_endpoint(member("a"));
889 bus.register_endpoint(member("b"));
890
891 let err = bus
892 .send(SignalFrame {
893 authenticated_peer: member("stranger"),
894 from: member("stranger"),
895 to: member("a"),
896 messages: vec![liveness("stranger", "a")],
897 })
898 .expect_err("unadmitted sender must be rejected");
899
900 assert_eq!(
901 err,
902 SignalTransportError::SenderNotAdmitted(member("stranger"))
903 );
904 assert_eq!(bus.drain(member("a")).len(), 0);
905 assert_eq!(bus.metrics().rejected_frames_total, 1);
906 }
907
908 #[test]
909 fn three_node_gossip_converges_over_secured_transport_with_one_member_down() {
910 let transport = SharedSignalTransport::default();
911 let limits = SignalPlaneLimits {
912 fanout: 1,
913 max_payload_messages: 2,
914 };
915 let members = vec![member("a"), member("b"), member("c")];
916
917 let mut node_a = TransportSignalPlane::new(member("a"), transport.clone(), limits);
918 let mut node_b = TransportSignalPlane::new(member("b"), transport.clone(), limits);
919 let mut node_c = TransportSignalPlane::new(member("c"), transport.clone(), limits);
920
921 node_a.set_members(members.clone());
922 node_b.set_members(members.clone());
923 node_c.set_members(members);
924 transport.unregister_endpoint(&member("c"));
925
926 let observation = SignalPlaneMessage::LivenessObservation(LivenessObservation {
927 observer: member("a"),
928 observed: member("c"),
929 incarnation: 2,
930 status: LivenessStatus::Unreachable,
931 });
932 node_a.publish(member("a"), observation.clone());
933
934 let bound = 6;
935 for _ in 0..bound {
936 node_a.advance_round();
937 node_b.advance_round();
938 node_c.advance_round();
939 if node_a.knows(&observation) && node_b.knows(&observation) {
940 break;
941 }
942 }
943
944 assert!(node_a.knows(&observation));
945 assert!(node_b.knows(&observation));
946 assert!(!node_c.knows(&observation));
947
948 let metrics = node_a.metrics();
949 assert!(metrics.sent_frames_total > 0);
950 assert!(metrics.sent_messages_total > 0);
951 assert!(metrics.fanout_cap >= 1);
952 assert!(metrics.payload_cap >= 2);
953 }
954}