Skip to main content

memberlist_core/delegate/
composite.rs

1use super::*;
2
3/// `CompositeDelegate` is a helpful struct to split the [`Delegate`] into multiple small delegates,
4/// so that users do not need to implement full [`Delegate`] when they only want to custom some methods
5/// in the [`Delegate`].
6pub struct CompositeDelegate<
7  I,
8  Address,
9  A = VoidDelegate<I, Address>,
10  C = VoidDelegate<I, Address>,
11  E = VoidDelegate<I, Address>,
12  M = VoidDelegate<I, Address>,
13  N = VoidDelegate<I, Address>,
14  P = VoidDelegate<I, Address>,
15> {
16  alive_delegate: A,
17  conflict_delegate: C,
18  event_delegate: E,
19  merge_delegate: M,
20  node_delegate: N,
21  ping_delegate: P,
22  _m: std::marker::PhantomData<(I, Address)>,
23}
24
25impl<I, Address> Default for CompositeDelegate<I, Address> {
26  fn default() -> Self {
27    Self::new()
28  }
29}
30
31impl<I, Address> CompositeDelegate<I, Address> {
32  /// Create a new `CompositeDelegate`
33  #[inline]
34  pub const fn new() -> Self {
35    Self {
36      alive_delegate: VoidDelegate::new(),
37      conflict_delegate: VoidDelegate::new(),
38      event_delegate: VoidDelegate::new(),
39      merge_delegate: VoidDelegate::new(),
40      node_delegate: VoidDelegate::new(),
41      ping_delegate: VoidDelegate::new(),
42      _m: std::marker::PhantomData,
43    }
44  }
45}
46
47impl<I, Address, A, C, E, M, N, P> CompositeDelegate<I, Address, A, C, E, M, N, P> {
48  /// Set the alive delegate
49  pub fn with_alive_delegate<NA>(
50    self,
51    alive_delegate: NA,
52  ) -> CompositeDelegate<I, Address, NA, C, E, M, N, P> {
53    CompositeDelegate {
54      alive_delegate,
55      conflict_delegate: self.conflict_delegate,
56      event_delegate: self.event_delegate,
57      merge_delegate: self.merge_delegate,
58      node_delegate: self.node_delegate,
59      ping_delegate: self.ping_delegate,
60      _m: std::marker::PhantomData,
61    }
62  }
63}
64
65impl<I, Address, A, C, E, M, N, P> CompositeDelegate<I, Address, A, C, E, M, N, P> {
66  /// Set the conflict delegate
67  pub fn with_conflict_delegate<NC>(
68    self,
69    conflict_delegate: NC,
70  ) -> CompositeDelegate<I, Address, A, NC, E, M, N, P> {
71    CompositeDelegate {
72      alive_delegate: self.alive_delegate,
73      conflict_delegate,
74      event_delegate: self.event_delegate,
75      merge_delegate: self.merge_delegate,
76      node_delegate: self.node_delegate,
77      ping_delegate: self.ping_delegate,
78      _m: std::marker::PhantomData,
79    }
80  }
81}
82
83impl<I, Address, A, C, E, M, N, P> CompositeDelegate<I, Address, A, C, E, M, N, P> {
84  /// Set the event delegate
85  pub fn with_event_delegate<NE>(
86    self,
87    event_delegate: NE,
88  ) -> CompositeDelegate<I, Address, A, C, NE, M, N, P> {
89    CompositeDelegate {
90      alive_delegate: self.alive_delegate,
91      conflict_delegate: self.conflict_delegate,
92      event_delegate,
93      merge_delegate: self.merge_delegate,
94      node_delegate: self.node_delegate,
95      ping_delegate: self.ping_delegate,
96      _m: std::marker::PhantomData,
97    }
98  }
99}
100
101impl<I, Address, A, C, E, M, N, P> CompositeDelegate<I, Address, A, C, E, M, N, P> {
102  /// Set the merge delegate
103  pub fn with_merge_delegate<NM>(
104    self,
105    merge_delegate: NM,
106  ) -> CompositeDelegate<I, Address, A, C, E, NM, N, P> {
107    CompositeDelegate {
108      alive_delegate: self.alive_delegate,
109      conflict_delegate: self.conflict_delegate,
110      event_delegate: self.event_delegate,
111      merge_delegate,
112      node_delegate: self.node_delegate,
113      ping_delegate: self.ping_delegate,
114      _m: std::marker::PhantomData,
115    }
116  }
117}
118
119impl<I, Address, A, C, E, M, N, P> CompositeDelegate<I, Address, A, C, E, M, N, P> {
120  /// Set the node delegate
121  pub fn with_node_delegate<NN>(
122    self,
123    node_delegate: NN,
124  ) -> CompositeDelegate<I, Address, A, C, E, M, NN, P> {
125    CompositeDelegate {
126      alive_delegate: self.alive_delegate,
127      conflict_delegate: self.conflict_delegate,
128      event_delegate: self.event_delegate,
129      merge_delegate: self.merge_delegate,
130      node_delegate,
131      ping_delegate: self.ping_delegate,
132      _m: std::marker::PhantomData,
133    }
134  }
135}
136
137impl<I, Address, A, C, E, M, N, P> CompositeDelegate<I, Address, A, C, E, M, N, P> {
138  /// Set the ping delegate
139  pub fn with_ping_delegate<NP>(
140    self,
141    ping_delegate: NP,
142  ) -> CompositeDelegate<I, Address, A, C, E, M, N, NP> {
143    CompositeDelegate {
144      alive_delegate: self.alive_delegate,
145      conflict_delegate: self.conflict_delegate,
146      event_delegate: self.event_delegate,
147      merge_delegate: self.merge_delegate,
148      node_delegate: self.node_delegate,
149      ping_delegate,
150      _m: std::marker::PhantomData,
151    }
152  }
153}
154
155#[cfg(any(feature = "test", test))]
156impl<I, Address, A, C, E, M, N, P> CompositeDelegate<I, Address, A, C, E, M, N, P> {
157  pub(crate) fn node_delegate(&self) -> &N {
158    &self.node_delegate
159  }
160
161  pub(crate) fn event_delegate(&self) -> &E {
162    &self.event_delegate
163  }
164
165  pub(crate) fn merge_delegate(&self) -> &M {
166    &self.merge_delegate
167  }
168
169  pub(crate) fn alive_delegate(&self) -> &A {
170    &self.alive_delegate
171  }
172
173  pub(crate) fn conflict_delegate(&self) -> &C {
174    &self.conflict_delegate
175  }
176
177  pub(crate) fn ping_delegate(&self) -> &P {
178    &self.ping_delegate
179  }
180}
181
182impl<I, Address, A, C, E, M, N, P> AliveDelegate for CompositeDelegate<I, Address, A, C, E, M, N, P>
183where
184  I: Id + Send + Sync + 'static,
185  Address: CheapClone + Send + Sync + 'static,
186  A: AliveDelegate<Id = I, Address = Address>,
187  C: ConflictDelegate<Id = I, Address = Address>,
188  E: EventDelegate<Id = I, Address = Address>,
189  M: MergeDelegate<Id = I, Address = Address>,
190  N: NodeDelegate,
191  P: PingDelegate<Id = I, Address = Address>,
192{
193  type Error = A::Error;
194  type Id = I;
195  type Address = Address;
196
197  async fn notify_alive(
198    &self,
199    peer: Arc<NodeState<Self::Id, Self::Address>>,
200  ) -> Result<(), Self::Error> {
201    self.alive_delegate.notify_alive(peer).await
202  }
203}
204
205impl<I, Address, A, C, E, M, N, P> MergeDelegate for CompositeDelegate<I, Address, A, C, E, M, N, P>
206where
207  I: Id + Send + Sync + 'static,
208  Address: CheapClone + Send + Sync + 'static,
209  A: AliveDelegate<Id = I, Address = Address>,
210  C: ConflictDelegate<Id = I, Address = Address>,
211  E: EventDelegate<Id = I, Address = Address>,
212  M: MergeDelegate<Id = I, Address = Address>,
213  N: NodeDelegate,
214  P: PingDelegate<Id = I, Address = Address>,
215{
216  type Error = M::Error;
217  type Id = I;
218  type Address = Address;
219
220  async fn notify_merge(
221    &self,
222    peers: Arc<[NodeState<Self::Id, Self::Address>]>,
223  ) -> Result<(), Self::Error> {
224    self.merge_delegate.notify_merge(peers).await
225  }
226}
227
228impl<I, Address, A, C, E, M, N, P> ConflictDelegate
229  for CompositeDelegate<I, Address, A, C, E, M, N, P>
230where
231  I: Id + Send + Sync + 'static,
232  Address: CheapClone + Send + Sync + 'static,
233  A: AliveDelegate<Id = I, Address = Address>,
234  C: ConflictDelegate<Id = I, Address = Address>,
235  E: EventDelegate<Id = I, Address = Address>,
236  M: MergeDelegate<Id = I, Address = Address>,
237  N: NodeDelegate,
238  P: PingDelegate<Id = I, Address = Address>,
239{
240  type Id = I;
241  type Address = Address;
242
243  async fn notify_conflict(
244    &self,
245    existing: Arc<NodeState<Self::Id, Self::Address>>,
246    other: Arc<NodeState<Self::Id, Self::Address>>,
247  ) {
248    self
249      .conflict_delegate
250      .notify_conflict(existing, other)
251      .await
252  }
253}
254
255impl<I, Address, A, C, E, M, N, P> PingDelegate for CompositeDelegate<I, Address, A, C, E, M, N, P>
256where
257  I: Id + Send + Sync + 'static,
258  Address: CheapClone + Send + Sync + 'static,
259  A: AliveDelegate<Id = I, Address = Address>,
260  C: ConflictDelegate<Id = I, Address = Address>,
261  E: EventDelegate<Id = I, Address = Address>,
262  M: MergeDelegate<Id = I, Address = Address>,
263  N: NodeDelegate,
264  P: PingDelegate<Id = I, Address = Address>,
265{
266  type Id = I;
267  type Address = Address;
268
269  async fn ack_payload(&self) -> Bytes {
270    self.ping_delegate.ack_payload().await
271  }
272
273  async fn notify_ping_complete(
274    &self,
275    node: Arc<NodeState<Self::Id, Self::Address>>,
276    rtt: std::time::Duration,
277    payload: Bytes,
278  ) {
279    self
280      .ping_delegate
281      .notify_ping_complete(node, rtt, payload)
282      .await
283  }
284
285  fn disable_reliable_pings(&self, target: &Self::Id) -> bool {
286    self.ping_delegate.disable_reliable_pings(target)
287  }
288}
289
290impl<I, Address, A, C, E, M, N, P> EventDelegate for CompositeDelegate<I, Address, A, C, E, M, N, P>
291where
292  I: Id + Send + Sync + 'static,
293  Address: CheapClone + Send + Sync + 'static,
294  A: AliveDelegate<Id = I, Address = Address>,
295  C: ConflictDelegate<Id = I, Address = Address>,
296  E: EventDelegate<Id = I, Address = Address>,
297  M: MergeDelegate<Id = I, Address = Address>,
298  N: NodeDelegate,
299  P: PingDelegate<Id = I, Address = Address>,
300{
301  type Id = I;
302
303  type Address = Address;
304
305  async fn notify_join(&self, node: Arc<NodeState<Self::Id, Self::Address>>) {
306    self.event_delegate.notify_join(node).await
307  }
308
309  async fn notify_leave(&self, node: Arc<NodeState<Self::Id, Self::Address>>) {
310    self.event_delegate.notify_leave(node).await
311  }
312
313  async fn notify_update(&self, node: Arc<NodeState<Self::Id, Self::Address>>) {
314    self.event_delegate.notify_update(node).await
315  }
316}
317
318impl<I, Address, A, C, E, M, N, P> NodeDelegate for CompositeDelegate<I, Address, A, C, E, M, N, P>
319where
320  I: Id + Send + Sync + 'static,
321  Address: CheapClone + Send + Sync + 'static,
322  A: AliveDelegate<Id = I, Address = Address>,
323  C: ConflictDelegate<Id = I, Address = Address>,
324  E: EventDelegate<Id = I, Address = Address>,
325  M: MergeDelegate<Id = I, Address = Address>,
326  N: NodeDelegate,
327  P: PingDelegate<Id = I, Address = Address>,
328{
329  async fn node_meta(&self, limit: usize) -> Meta {
330    self.node_delegate.node_meta(limit).await
331  }
332
333  async fn notify_message(&self, msg: Cow<'_, [u8]>) {
334    self.node_delegate.notify_message(msg).await
335  }
336
337  async fn broadcast_messages<F>(
338    &self,
339    limit: usize,
340    encoded_len: F,
341  ) -> impl Iterator<Item = Bytes> + Send
342  where
343    F: Fn(Bytes) -> (usize, Bytes) + Send + Sync + 'static,
344  {
345    self
346      .node_delegate
347      .broadcast_messages(limit, encoded_len)
348      .await
349  }
350
351  async fn local_state(&self, join: bool) -> Bytes {
352    self.node_delegate.local_state(join).await
353  }
354
355  async fn merge_remote_state(&self, buf: &[u8], join: bool) {
356    self.node_delegate.merge_remote_state(buf, join).await
357  }
358}
359
360impl<I, Address, A, C, E, M, N, P> Delegate for CompositeDelegate<I, Address, A, C, E, M, N, P>
361where
362  I: Id + Send + Sync + 'static,
363  Address: CheapClone + Send + Sync + 'static,
364  A: AliveDelegate<Id = I, Address = Address>,
365  C: ConflictDelegate<Id = I, Address = Address>,
366  E: EventDelegate<Id = I, Address = Address>,
367  M: MergeDelegate<Id = I, Address = Address>,
368  N: NodeDelegate,
369  P: PingDelegate<Id = I, Address = Address>,
370{
371  type Address = Address;
372  type Id = I;
373}