Skip to main content

reddb_server/replication/
signal_plane.rs

1//! Signal-plane dissemination seam (issue #1842, parent #1832, ADR 0073).
2//!
3//! This module names the non-authoritative peer-to-peer signal layer used by
4//! admitted cluster members to share routing and health hints. The vocabulary is
5//! closed over observations only: liveness, health inputs, load samples,
6//! catalog-version hints, and topology hints. There is no membership admission,
7//! ownership transition, vote, or bootstrap-state variant, so those authority
8//! facts cannot be represented by this seam.
9
10use std::collections::{BTreeMap, BTreeSet, VecDeque};
11use std::sync::{Arc, Mutex};
12
13use super::MemberId;
14
15/// Peer reachability state observed by an admitted member.
16#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
17pub enum LivenessStatus {
18    Alive,
19    Suspect,
20    Unreachable,
21}
22
23/// Coarse load bucket used by signal-plane load samples.
24#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
25pub enum LoadBucket {
26    Low,
27    Medium,
28    High,
29    Saturated,
30}
31
32/// A peer reachability observation. This is a hint, not membership authority.
33#[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/// Bounded health inputs that may influence local scoring or refresh timing.
42#[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/// Coarse capacity and throughput sample for a known member.
52#[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/// Non-authoritative catalog and topology generation hint.
63#[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/// Routing-adjacent topology hint for an already-known member.
72#[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/// Closed signal-plane vocabulary.
81///
82/// Deliberately absent: membership admission, ownership transitions, votes, and
83/// bootstrap state. Adding any authority-bearing variant is a boundary change,
84/// not an implementation detail.
85#[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/// A signal delivered to a local consumer.
95#[derive(Debug, Clone, PartialEq, Eq)]
96pub struct ReceivedSignal {
97    pub from: MemberId,
98    pub to: MemberId,
99    pub message: SignalPlaneMessage,
100}
101
102/// Narrow internal seam for signal dissemination.
103///
104/// Implementations receive authoritative membership as input and disseminate the
105/// closed signal vocabulary. They do not expose join, ownership, voting, or
106/// bootstrap APIs.
107pub trait SignalPlane {
108    /// Replace the admitted-member set this signal plane may talk to.
109    fn set_members(&mut self, members: Vec<MemberId>);
110
111    /// Publish a non-authoritative signal from an admitted member.
112    fn publish(&mut self, from: MemberId, message: SignalPlaneMessage);
113
114    /// Drain signals delivered to `member`.
115    fn drain_received(&mut self, member: &MemberId) -> Vec<ReceivedSignal>;
116}
117
118/// Per-round bounds for signal-plane gossip traffic.
119#[derive(Debug, Clone, Copy, PartialEq, Eq)]
120pub struct SignalPlaneLimits {
121    /// Maximum peers sampled by one member in a round.
122    pub fanout: usize,
123    /// Maximum closed-vocabulary messages carried in one frame.
124    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/// Signal-plane traffic counters exported by the transport and engines.
137#[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/// One authenticated intra-cluster signal frame.
151#[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/// Rejection from the secured intra-cluster signal transport.
160#[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/// Secured member-to-member signal transport for admitted cluster members.
176#[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/// Cloneable handle for the secured intra-cluster signal transport.
276#[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/// SWIM-style signal-plane engine over the secured intra-cluster transport.
333#[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/// Deterministic fault schedule for the in-process simulated-peer engine.
498#[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    /// Drop all sends on one directed link during `round`.
507    pub fn drop_delivery(&mut self, round: u64, from: MemberId, to: MemberId) {
508        self.dropped.insert(ScheduledLink { round, from, to });
509    }
510
511    /// Add `extra_rounds` of delay to one directed link during `round`.
512    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    /// Partition two members from each other for an inclusive round range.
518    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/// Deterministic in-process signal plane for seam tests and future consumers.
560#[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    /// Advance one deterministic gossip round.
584    pub fn advance_round(&mut self) {
585        self.queue_known_signals();
586        self.deliver_due_signals();
587        self.round += 1;
588    }
589
590    /// Members whose local simulated peer has learned `message`.
591    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}