Skip to main content

hyperlane_plugin_websocket/
impl.rs

1use super::*;
2
3/// Allows `String` to be used as a broadcast identifier.
4impl BroadcastTypeTrait for String {}
5
6/// Allows string slices to be used as broadcast identifiers.
7impl BroadcastTypeTrait for &str {}
8
9/// Allows `char` to be used as a broadcast identifier.
10impl BroadcastTypeTrait for char {}
11
12/// Allows `bool` to be used as a broadcast identifier.
13impl BroadcastTypeTrait for bool {}
14
15/// Allows `i8` to be used as a broadcast identifier.
16impl BroadcastTypeTrait for i8 {}
17
18/// Allows `i16` to be used as a broadcast identifier.
19impl BroadcastTypeTrait for i16 {}
20
21/// Allows `i32` to be used as a broadcast identifier.
22impl BroadcastTypeTrait for i32 {}
23
24/// Allows `i64` to be used as a broadcast identifier.
25impl BroadcastTypeTrait for i64 {}
26
27/// Allows `i128` to be used as a broadcast identifier.
28impl BroadcastTypeTrait for i128 {}
29
30/// Allows `isize` to be used as a broadcast identifier.
31impl BroadcastTypeTrait for isize {}
32
33/// Allows `u8` to be used as a broadcast identifier.
34impl BroadcastTypeTrait for u8 {}
35
36/// Allows `u16` to be used as a broadcast identifier.
37impl BroadcastTypeTrait for u16 {}
38
39/// Allows `u32` to be used as a broadcast identifier.
40impl BroadcastTypeTrait for u32 {}
41
42/// Allows `u64` to be used as a broadcast identifier.
43impl BroadcastTypeTrait for u64 {}
44
45/// Allows `u128` to be used as a broadcast identifier.
46impl BroadcastTypeTrait for u128 {}
47
48/// Allows `usize` to be used as a broadcast identifier.
49impl BroadcastTypeTrait for usize {}
50
51/// Allows `f32` to be used as a broadcast identifier.
52impl BroadcastTypeTrait for f32 {}
53
54/// Allows `f64` to be used as a broadcast identifier.
55impl BroadcastTypeTrait for f64 {}
56
57/// Allows `IpAddr` to be used as a broadcast identifier.
58impl BroadcastTypeTrait for IpAddr {}
59
60/// Allows `Ipv4Addr` to be used as a broadcast identifier.
61impl BroadcastTypeTrait for Ipv4Addr {}
62
63/// Allows `Ipv6Addr` to be used as a broadcast identifier.
64impl BroadcastTypeTrait for Ipv6Addr {}
65
66/// Allows `SocketAddr` to be used as a broadcast identifier.
67impl BroadcastTypeTrait for SocketAddr {}
68
69/// Allows `NonZeroU8` to be used as a broadcast identifier.
70impl BroadcastTypeTrait for NonZeroU8 {}
71
72/// Allows `NonZeroU16` to be used as a broadcast identifier.
73impl BroadcastTypeTrait for NonZeroU16 {}
74
75/// Allows `NonZeroU32` to be used as a broadcast identifier.
76impl BroadcastTypeTrait for NonZeroU32 {}
77
78/// Allows `NonZeroU64` to be used as a broadcast identifier.
79impl BroadcastTypeTrait for NonZeroU64 {}
80
81/// Allows `NonZeroU128` to be used as a broadcast identifier.
82impl BroadcastTypeTrait for NonZeroU128 {}
83
84/// Allows `NonZeroUsize` to be used as a broadcast identifier.
85impl BroadcastTypeTrait for NonZeroUsize {}
86
87/// Allows `NonZeroI8` to be used as a broadcast identifier.
88impl BroadcastTypeTrait for NonZeroI8 {}
89
90/// Allows `NonZeroI16` to be used as a broadcast identifier.
91impl BroadcastTypeTrait for NonZeroI16 {}
92
93/// Allows `NonZeroI32` to be used as a broadcast identifier.
94impl BroadcastTypeTrait for NonZeroI32 {}
95
96/// Allows `NonZeroI64` to be used as a broadcast identifier.
97impl BroadcastTypeTrait for NonZeroI64 {}
98
99/// Allows `NonZeroI128` to be used as a broadcast identifier.
100impl BroadcastTypeTrait for NonZeroI128 {}
101
102/// Allows `NonZeroIsize` to be used as a broadcast identifier.
103impl BroadcastTypeTrait for NonZeroIsize {}
104
105/// Allows `Infallible` to be used as a broadcast identifier.
106impl BroadcastTypeTrait for Infallible {}
107
108/// Allows references to `String` to be used as broadcast identifiers.
109impl BroadcastTypeTrait for &String {}
110
111/// Allows double references to string slices to be used as broadcast identifiers.
112impl BroadcastTypeTrait for &&str {}
113
114/// Allows references to `char` to be used as broadcast identifiers.
115impl BroadcastTypeTrait for &char {}
116
117/// Allows references to `bool` to be used as broadcast identifiers.
118impl BroadcastTypeTrait for &bool {}
119
120/// Allows references to `i8` to be used as broadcast identifiers.
121impl BroadcastTypeTrait for &i8 {}
122
123/// Allows references to `i16` to be used as broadcast identifiers.
124impl BroadcastTypeTrait for &i16 {}
125
126/// Allows references to `i32` to be used as broadcast identifiers.
127impl BroadcastTypeTrait for &i32 {}
128
129/// Allows references to `i64` to be used as broadcast identifiers.
130impl BroadcastTypeTrait for &i64 {}
131
132/// Allows references to `i128` to be used as broadcast identifiers.
133impl BroadcastTypeTrait for &i128 {}
134
135/// Allows references to `isize` to be used as broadcast identifiers.
136impl BroadcastTypeTrait for &isize {}
137
138/// Allows references to `u8` to be used as broadcast identifiers.
139impl BroadcastTypeTrait for &u8 {}
140
141/// Allows references to `u16` to be used as broadcast identifiers.
142impl BroadcastTypeTrait for &u16 {}
143
144/// Allows references to `u32` to be used as broadcast identifiers.
145impl BroadcastTypeTrait for &u32 {}
146
147/// Implements `BroadcastTypeTrait` for `&u64`.
148///
149/// This allows references to `u64` to be used as a broadcast identifier.
150impl BroadcastTypeTrait for &u64 {}
151
152/// Implements `BroadcastTypeTrait` for `&u128`.
153///
154/// This allows references to `u128` to be used as a broadcast identifier.
155impl BroadcastTypeTrait for &u128 {}
156
157/// Implements `BroadcastTypeTrait` for `&usize`.
158///
159/// This allows references to `usize` to be used as a broadcast identifier.
160impl BroadcastTypeTrait for &usize {}
161
162/// Implements `BroadcastTypeTrait` for `&f32`.
163///
164/// This allows references to `f32` to be used as a broadcast identifier.
165impl BroadcastTypeTrait for &f32 {}
166
167/// Implements `BroadcastTypeTrait` for `&f64`.
168///
169/// This allows references to `f64` to be used as a broadcast identifier.
170impl BroadcastTypeTrait for &f64 {}
171
172/// Implements `BroadcastTypeTrait` for `&IpAddr`.
173///
174/// This allows references to `IpAddr` to be used as a broadcast identifier.
175impl BroadcastTypeTrait for &IpAddr {}
176
177/// Implements `BroadcastTypeTrait` for `&Ipv4Addr`.
178///
179/// This allows references to `Ipv4Addr` to be used as a broadcast identifier.
180impl BroadcastTypeTrait for &Ipv4Addr {}
181
182/// Implements `BroadcastTypeTrait` for `&Ipv6Addr`.
183///
184/// This allows references to `Ipv6Addr` to be used as a broadcast identifier.
185impl BroadcastTypeTrait for &Ipv6Addr {}
186
187/// Implements `BroadcastTypeTrait` for `&SocketAddr`.
188///
189/// This allows references to `SocketAddr` to be used as a broadcast identifier.
190impl BroadcastTypeTrait for &SocketAddr {}
191
192/// Implements `BroadcastTypeTrait` for `&NonZeroU8`.
193///
194/// This allows references to `NonZeroU8` to be used as a broadcast identifier.
195impl BroadcastTypeTrait for &NonZeroU8 {}
196
197/// Implements `BroadcastTypeTrait` for `&NonZeroU16`.
198///
199/// This allows references to `NonZeroU16` to be used as a broadcast identifier.
200impl BroadcastTypeTrait for &NonZeroU16 {}
201
202/// Implements `BroadcastTypeTrait` for `&NonZeroU32`.
203///
204/// This allows references to `NonZeroU32` to be used as a broadcast identifier.
205impl BroadcastTypeTrait for &NonZeroU32 {}
206
207/// Implements `BroadcastTypeTrait` for `&NonZeroU64`.
208///
209/// This allows references to `NonZeroU64` to be used as a broadcast identifier.
210impl BroadcastTypeTrait for &NonZeroU64 {}
211
212/// Implements `BroadcastTypeTrait` for `&NonZeroU128`.
213///
214/// This allows references to `NonZeroU128` to be used as a broadcast identifier.
215impl BroadcastTypeTrait for &NonZeroU128 {}
216
217/// Implements `BroadcastTypeTrait` for `&NonZeroUsize`.
218///
219/// This allows references to `NonZeroUsize` to be used as a broadcast identifier.
220impl BroadcastTypeTrait for &NonZeroUsize {}
221
222/// Implements `BroadcastTypeTrait` for `&NonZeroI8`.
223///
224/// This allows references to `NonZeroI8` to be used as a broadcast identifier.
225impl BroadcastTypeTrait for &NonZeroI8 {}
226
227/// Implements `BroadcastTypeTrait` for `&NonZeroI16`.
228///
229/// This allows references to `NonZeroI16` to be used as a broadcast identifier.
230impl BroadcastTypeTrait for &NonZeroI16 {}
231
232/// Implements `BroadcastTypeTrait` for `&NonZeroI32`.
233///
234/// This allows references to `NonZeroI32` to be used as a broadcast identifier.
235impl BroadcastTypeTrait for &NonZeroI32 {}
236
237/// Implements `BroadcastTypeTrait` for `&NonZeroI64`.
238///
239/// This allows references to `NonZeroI64` to be used as a broadcast identifier.
240impl BroadcastTypeTrait for &NonZeroI64 {}
241
242/// Implements `BroadcastTypeTrait` for `&NonZeroI128`.
243///
244/// This allows references to `NonZeroI128` to be used as a broadcast identifier.
245impl BroadcastTypeTrait for &NonZeroI128 {}
246
247/// Implements `BroadcastTypeTrait` for `&NonZeroIsize`.
248///
249/// This allows references to `NonZeroIsize` to be used as a broadcast identifier.
250impl BroadcastTypeTrait for &NonZeroIsize {}
251
252/// Implements `BroadcastTypeTrait` for `&Infallible`.
253///
254/// This allows references to `Infallible` to be used as a broadcast identifier.
255impl BroadcastTypeTrait for &Infallible {}
256
257/// Implements the `Default` trait for `BroadcastType`.
258///
259/// The default value is `BroadcastType::Unknown`.
260///
261/// # Type Parameters
262///
263/// - `BroadcastTypeTrait`: The type parameter for `BroadcastType`, which must implement `BroadcastTypeTrait`.
264impl<B> Default for BroadcastType<B>
265where
266    B: BroadcastTypeTrait,
267{
268    /// Returns the default `BroadcastType`, which is `BroadcastType::Unknown`.
269    #[inline(always)]
270    fn default() -> Self {
271        BroadcastType::Unknown
272    }
273}
274
275impl<B> BroadcastType<B>
276where
277    B: BroadcastTypeTrait,
278{
279    /// Generates a unique key string for a given broadcast type.
280    ///
281    /// For point-to-point types, the keys are sorted to ensure consistent key generation
282    /// regardless of the order of the input keys.
283    ///
284    /// # Arguments
285    ///
286    /// - `BroadcastType<B>` - The broadcast type for which to generate the key.
287    ///
288    /// # Returns
289    ///
290    /// - `String` - The unique key string for the broadcast type.
291    #[inline(always)]
292    pub fn get_key(broadcast_type: BroadcastType<B>) -> String {
293        match broadcast_type {
294            BroadcastType::PointToPoint(key1, key2) => {
295                let (first_key, second_key) = if key1 <= key2 {
296                    (key1, key2)
297                } else {
298                    (key2, key1)
299                };
300                format!(
301                    "{}-{}-{}",
302                    POINT_TO_POINT_KEY,
303                    first_key.to_string(),
304                    second_key.to_string()
305                )
306            }
307            BroadcastType::PointToGroup(key) => {
308                format!("{}-{}", POINT_TO_GROUP_KEY, key.to_string())
309            }
310            BroadcastType::Unknown => String::new(),
311        }
312    }
313}
314
315impl<'a, B> WebSocketConfig<'a, B>
316where
317    B: BroadcastTypeTrait,
318{
319    /// Creates a new WebSocket configuration with the given context.
320    ///
321    /// # Arguments
322    ///
323    /// - `&'a mut Stream` - The stream object serving this WebSocket.
324    /// - `&'a mut Context` - The context object to associate with the WebSocket.
325    ///
326    /// # Returns
327    ///
328    /// - `WebSocketConfig<B>` - A new WebSocket configuration instance.
329    #[inline(always)]
330    pub fn new(stream: &'a mut Stream, context: &'a mut Context) -> Self {
331        Self {
332            stream,
333            context,
334            capacity: DEFAULT_BROADCAST_SENDER_CAPACITY,
335            broadcast_type: BroadcastType::default(),
336            connected_hook: Hook::default_handler(),
337            request_hook: Hook::default_handler(),
338            sended_hook: Hook::default_handler(),
339            closed_hook: Hook::default_handler(),
340        }
341    }
342}
343
344impl<'a, B> WebSocketConfig<'a, B>
345where
346    B: BroadcastTypeTrait,
347{
348    /// Sets the capacity for the broadcast sender.
349    ///
350    /// # Arguments
351    ///
352    /// - `Capacity` - The desired capacity.
353    ///
354    /// # Returns
355    ///
356    /// - `WebSocketConfig<B>` - The modified WebSocket configuration instance.
357    #[inline(always)]
358    pub fn set_capacity(mut self, capacity: Capacity) -> Self {
359        self.capacity = capacity;
360        self
361    }
362
363    /// Sets the context for the WebSocket connection.
364    ///
365    /// # Arguments
366    ///
367    /// - `&'a mut Context` - The context object to associate with the WebSocket.
368    ///
369    /// # Returns
370    ///
371    /// - `WebSocketConfig<B>` - The modified WebSocket configuration instance.
372    #[inline(always)]
373    pub fn set_context(mut self, context: &'a mut Context) -> Self {
374        self.context = context;
375        self
376    }
377
378    /// Sets the broadcast type for the WebSocket connection.
379    ///
380    /// # Arguments
381    ///
382    /// - `BroadcastType<B>` - The broadcast type to use for this WebSocket.
383    ///
384    /// # Returns
385    ///
386    /// - `WebSocketConfig<B>` - The modified WebSocket configuration instance.
387    #[inline(always)]
388    pub fn set_broadcast_type(mut self, broadcast_type: BroadcastType<B>) -> Self {
389        self.broadcast_type = broadcast_type;
390        self
391    }
392
393    /// Returns a mutable reference to the stream served by this configuration.
394    ///
395    /// # Returns
396    ///
397    /// - `&mut Stream` - A mutable reference to the stream.
398    #[inline(always)]
399    pub fn get_stream(&mut self) -> &mut Stream {
400        self.stream
401    }
402
403    /// Retrieves a reference to the context associated with this configuration.
404    ///
405    /// # Returns
406    ///
407    /// - `&mut Context` - A reference to the context object.
408    #[inline(always)]
409    pub fn get_context(&mut self) -> &mut Context {
410        self.context
411    }
412
413    /// Retrieves the capacity configured for the broadcast sender.
414    ///
415    /// # Returns
416    ///
417    /// - `Capacity` - The capacity.
418    #[inline(always)]
419    pub fn get_capacity(&self) -> Capacity {
420        self.capacity
421    }
422
423    /// Retrieves a reference to the broadcast type configured for this WebSocket.
424    ///
425    /// # Returns
426    ///
427    /// - `&BroadcastType<B>` - A reference to the broadcast type object.
428    #[inline(always)]
429    pub fn get_broadcast_type(&self) -> &BroadcastType<B> {
430        &self.broadcast_type
431    }
432
433    /// Sets the connected hook handler.
434    ///
435    /// This hook is executed when the WebSocket connection is established.
436    ///
437    /// # Type Parameters
438    ///
439    /// - `S`: The hook type, which must implement `ServerHook`.
440    ///
441    /// # Returns
442    ///
443    /// The modified `WebSocketConfig` instance.
444    ///
445    /// # Examples
446    ///
447    /// ```rust,ignore
448    /// struct MyConnectedHook;
449    /// impl ServerHook for MyConnectedHook {
450    ///     async fn new(_ctx: &Context) -> Self { Self }
451    ///     async fn handle(self, ctx: &Context) { /* ... */ }
452    /// }
453    ///
454    /// let config = WebSocketConfig::new()
455    ///     .set_connected_hook::<MyConnectedHook>();
456    /// ```
457    #[inline(always)]
458    pub fn set_connected_hook<S>(mut self) -> Self
459    where
460        S: ServerHook,
461    {
462        self.connected_hook = Hook::factory::<S>();
463        self
464    }
465
466    /// Sets the request hook handler.
467    ///
468    /// This hook is executed when a new request is received on the WebSocket.
469    ///
470    /// # Type Parameters
471    ///
472    /// - `S`: The hook type, which must implement `ServerHook`.
473    ///
474    /// # Returns
475    ///
476    /// The modified `WebSocketConfig` instance.
477    ///
478    /// # Examples
479    ///
480    /// ```rust,ignore
481    /// struct MyRequestHook;
482    /// impl ServerHook for MyRequestHook {
483    ///     async fn new(_ctx: &Context) -> Self { Self }
484    ///     async fn handle(self, ctx: &Context) { /* ... */ }
485    /// }
486    ///
487    /// let config = WebSocketConfig::new()
488    ///     .set_request_hook::<MyRequestHook>();
489    /// ```
490    #[inline(always)]
491    pub fn set_request_hook<S>(mut self) -> Self
492    where
493        S: ServerHook,
494    {
495        self.request_hook = Hook::factory::<S>();
496        self
497    }
498
499    /// Sets the sended hook handler.
500    ///
501    /// This hook is executed after a message has been successfully sent over the WebSocket.
502    ///
503    /// # Type Parameters
504    ///
505    /// - `S`: The hook type, which must implement `ServerHook`.
506    ///
507    /// # Returns
508    ///
509    /// The modified `WebSocketConfig` instance.
510    ///
511    /// # Examples
512    ///
513    /// ```rust,ignore
514    /// struct MySendedHook;
515    /// impl ServerHook for MySendedHook {
516    ///     async fn new(_ctx: &Context) -> Self { Self }
517    ///     async fn handle(self, ctx: &Context) { /* ... */ }
518    /// }
519    ///
520    /// let config = WebSocketConfig::new()
521    ///     .set_sended_hook::<MySendedHook>();
522    /// ```
523    #[inline(always)]
524    pub fn set_sended_hook<S>(mut self) -> Self
525    where
526        S: ServerHook,
527    {
528        self.sended_hook = Hook::factory::<S>();
529        self
530    }
531
532    /// Sets the closed hook handler.
533    ///
534    /// This hook is executed when the WebSocket connection is closed.
535    ///
536    /// # Type Parameters
537    ///
538    /// - `S`: The hook type, which must implement `ServerHook`.
539    ///
540    /// # Returns
541    ///
542    /// The modified `WebSocketConfig` instance.
543    ///
544    /// # Examples
545    ///
546    /// ```rust,ignore
547    /// struct MyClosedHook;
548    /// impl ServerHook for MyClosedHook {
549    ///     async fn new(_ctx: &Context) -> Self { Self }
550    ///     async fn handle(self, ctx: &Context) { /* ... */ }
551    /// }
552    ///
553    /// let config = WebSocketConfig::new()
554    ///     .set_closed_hook::<MyClosedHook>();
555    /// ```
556    #[inline(always)]
557    pub fn set_closed_hook<S>(mut self) -> Self
558    where
559        S: ServerHook,
560    {
561        self.closed_hook = Hook::factory::<S>();
562        self
563    }
564
565    /// Retrieves a reference to the connected hook handler.
566    ///
567    /// # Returns
568    ///
569    /// - `&ServerHookHandler` - A reference to the connected hook handler.
570    #[inline(always)]
571    pub fn get_connected_hook(&self) -> &ServerHookHandler {
572        &self.connected_hook
573    }
574
575    /// Retrieves a reference to the request hook handler.
576    ///
577    /// # Returns
578    ///
579    /// - `&ServerHookHandler` - A reference to the request hook handler.
580    #[inline(always)]
581    pub fn get_request_hook(&self) -> &ServerHookHandler {
582        &self.request_hook
583    }
584
585    /// Retrieves a reference to the sended hook handler.
586    ///
587    /// # Returns
588    ///
589    /// - `&ServerHookHandler` - A reference to the sended hook handler.
590    #[inline(always)]
591    pub fn get_sended_hook(&self) -> &ServerHookHandler {
592        &self.sended_hook
593    }
594
595    /// Retrieves a reference to the closed hook handler.
596    ///
597    /// # Returns
598    ///
599    /// - `&ServerHookHandler` - A reference to the closed hook handler.
600    #[inline(always)]
601    pub fn get_closed_hook(&self) -> &ServerHookHandler {
602        &self.closed_hook
603    }
604}
605
606impl WebSocket {
607    /// Creates a new WebSocket instance.
608    ///
609    /// Initializes with a default broadcast map.
610    ///
611    /// # Returns
612    ///
613    /// - `WebSocket` - A new WebSocket instance.
614    #[inline(always)]
615    pub fn new() -> Self {
616        Self::default()
617    }
618
619    /// Returns a shared reference to the internal broadcast map.
620    ///
621    /// Hand-written accessor: the `WebSocket` struct does not derive the lombok
622    /// `Data` macro, and ยง17.3 forbids reading `self.broadcast_map` directly.
623    ///
624    /// # Returns
625    ///
626    /// - `&BroadcastMap<Vec<u8>>` - A shared reference to the internal broadcast map.
627    #[inline(always)]
628    pub fn get_broadcast_map(&self) -> &BroadcastMap<Vec<u8>> {
629        &self.broadcast_map
630    }
631
632    /// Subscribes to a broadcast type or inserts a new one if it doesn't exist.
633    ///
634    /// # Type Parameters
635    ///
636    /// - `BroadcastTypeTrait`: The type implementing `BroadcastTypeTrait`.
637    ///
638    /// # Arguments
639    ///
640    /// - `BroadcastType<B>` - The broadcast type to subscribe to.
641    /// - `Capacity` - The capacity for the broadcast sender if a new one is inserted.
642    ///
643    /// # Returns
644    ///
645    /// - `BroadcastMapReceiver<Vec<u8>>` - A broadcast map receiver for the specified broadcast type.
646    #[inline(always)]
647    fn subscribe_unwrap_or_insert<B>(
648        &self,
649        broadcast_type: BroadcastType<B>,
650        capacity: Capacity,
651    ) -> BroadcastMapReceiver<Vec<u8>>
652    where
653        B: BroadcastTypeTrait,
654    {
655        let key: String = BroadcastType::get_key(broadcast_type);
656        self.get_broadcast_map().subscribe_or_insert(&key, capacity)
657    }
658
659    /// Subscribes to a point-to-point broadcast.
660    ///
661    /// # Type Parameters
662    ///
663    /// - `BroadcastTypeTrait`: The type implementing `BroadcastTypeTrait`.
664    ///
665    /// # Arguments
666    ///
667    /// - `&B` - The first identifier for the point-to-point communication.
668    /// - `&B` - The second identifier for the point-to-point communication.
669    /// - `Capacity` - The capacity for the broadcast sender.
670    ///
671    /// # Returns
672    ///
673    /// - `BroadcastMapReceiver<Vec<u8>>` - A broadcast map receiver for the point-to-point broadcast.
674    #[inline(always)]
675    fn point_to_point<B>(
676        &self,
677        key1: &B,
678        key2: &B,
679        capacity: Capacity,
680    ) -> BroadcastMapReceiver<Vec<u8>>
681    where
682        B: BroadcastTypeTrait,
683    {
684        self.subscribe_unwrap_or_insert(
685            BroadcastType::PointToPoint(key1.clone(), key2.clone()),
686            capacity,
687        )
688    }
689
690    /// Subscribes to a point-to-group broadcast.
691    ///
692    /// # Type Parameters
693    ///
694    /// - `BroadcastTypeTrait`: The type implementing `BroadcastTypeTrait`.
695    ///
696    /// # Arguments
697    ///
698    /// - `&B` - The identifier for the group.
699    /// - `Capacity` - The capacity for the broadcast sender.
700    ///
701    /// # Returns
702    ///
703    /// - `BroadcastMapReceiver<Vec<u8>>` - A broadcast map receiver for the point-to-group broadcast.
704    #[inline(always)]
705    fn point_to_group<B>(&self, key: &B, capacity: Capacity) -> BroadcastMapReceiver<Vec<u8>>
706    where
707        B: BroadcastTypeTrait,
708    {
709        self.subscribe_unwrap_or_insert(BroadcastType::PointToGroup(key.clone()), capacity)
710    }
711
712    /// Retrieves the current receiver count for a given broadcast type.
713    ///
714    /// # Type Parameters
715    ///
716    /// - `BroadcastTypeTrait`: The type implementing `BroadcastTypeTrait`.
717    ///
718    /// # Arguments
719    ///
720    /// - `BroadcastType<B>` - The broadcast type for which to get the receiver count.
721    ///
722    /// # Returns
723    ///
724    /// - `ReceiverCount` - The number of active receivers for the broadcast type, or 0 if not found.
725    #[inline(always)]
726    pub fn receiver_count<B>(&self, broadcast_type: BroadcastType<B>) -> ReceiverCount
727    where
728        B: BroadcastTypeTrait,
729    {
730        let key: String = BroadcastType::get_key(broadcast_type);
731        self.get_broadcast_map().receiver_count(&key).unwrap_or(0)
732    }
733
734    /// Calculates the receiver count before a connection is established.
735    ///
736    /// Ensures the count does not exceed the maximum allowed value minus one.
737    ///
738    /// # Type Parameters
739    ///
740    /// - `BroadcastTypeTrait`: The type implementing `BroadcastTypeTrait`.
741    ///
742    /// # Arguments
743    ///
744    /// - `BroadcastType<B>` - The broadcast type for which to get the receiver count.
745    ///
746    /// # Returns
747    ///
748    /// - `ReceiverCount` - The receiver count after the connection is established.
749    #[inline(always)]
750    pub fn receiver_count_before_connected<B>(
751        &self,
752        broadcast_type: BroadcastType<B>,
753    ) -> ReceiverCount
754    where
755        B: BroadcastTypeTrait,
756    {
757        let count: ReceiverCount = self.receiver_count(broadcast_type);
758        count.clamp(0, ReceiverCount::MAX - 1) + 1
759    }
760
761    /// Calculates the receiver count after a connection is closed.
762    ///
763    /// Ensures the count does not go below 0.
764    ///
765    /// # Type Parameters
766    ///
767    /// - `BroadcastTypeTrait`: The type implementing `BroadcastTypeTrait`.
768    ///
769    /// # Arguments
770    ///
771    /// - `BroadcastType<B>` - The broadcast type for which to get the receiver count.
772    ///
773    /// # Returns
774    ///
775    /// - `ReceiverCount` - The receiver count after the connection is closed.
776    #[inline(always)]
777    pub fn receiver_count_after_closed<B>(&self, broadcast_type: BroadcastType<B>) -> ReceiverCount
778    where
779        B: BroadcastTypeTrait,
780    {
781        let count: ReceiverCount = self.receiver_count(broadcast_type);
782        count.clamp(1, ReceiverCount::MAX) - 1
783    }
784
785    /// Sends data to all active receivers for a given broadcast type.
786    ///
787    /// # Type Parameters
788    ///
789    /// - `Into<Vec<u8>>`: The type of data to send, which must be convertible to `Vec<u8>`.
790    /// - `BroadcastTypeTrait`: The type implementing `BroadcastTypeTrait`.
791    ///
792    /// # Arguments
793    ///
794    /// - `BroadcastType<B>` - The broadcast type to which to send the data.
795    /// - `T` - The data to send.
796    ///
797    /// # Returns
798    ///
799    /// - `Result<Option<ReceiverCount>, SendError<Vec<u8>>>` - A result indicating the success or failure of the send operation.
800    #[inline(always)]
801    pub fn try_send<T, B>(
802        &self,
803        broadcast_type: BroadcastType<B>,
804        data: T,
805    ) -> Result<Option<ReceiverCount>, SendError<Vec<u8>>>
806    where
807        T: Into<Vec<u8>>,
808        B: BroadcastTypeTrait,
809    {
810        let key: String = BroadcastType::get_key(broadcast_type);
811        self.get_broadcast_map().try_send(&key, data.into())
812    }
813
814    /// Sends data to all active receivers for a given broadcast type.
815    ///
816    /// This method panics if the send operation fails.
817    ///
818    /// # Type Parameters
819    ///
820    /// - `Into<Vec<u8>>`: The type of data to send, which must be convertible to `Vec<u8>`.
821    /// - `BroadcastTypeTrait`: The type implementing `BroadcastTypeTrait`.
822    ///
823    /// # Arguments
824    ///
825    /// - `BroadcastType<B>` - The broadcast type to which to send the data.
826    /// - `T` - The data to send.
827    ///
828    /// # Returns
829    ///
830    /// - `Option<ReceiverCount>` - The receiver count if the send operation succeeds.
831    ///
832    /// # Panics
833    ///
834    /// Panics if the send operation fails.
835    #[inline(always)]
836    pub fn send<T, B>(&self, broadcast_type: BroadcastType<B>, data: T) -> Option<ReceiverCount>
837    where
838        T: Into<Vec<u8>>,
839        B: BroadcastTypeTrait,
840    {
841        self.try_send(broadcast_type, data).unwrap()
842    }
843
844    /// Runs the WebSocket connection, handling incoming requests and outgoing messages.
845    ///
846    /// This asynchronous function continuously monitors for new WebSocket requests
847    /// and incoming broadcast messages, processing them according to the configured hooks.
848    ///
849    /// # Type Parameters
850    ///
851    /// - `BroadcastTypeTrait`: The type implementing `BroadcastTypeTrait`.
852    ///
853    /// # Arguments
854    ///
855    /// - `WebSocketConfig<'_, B>` - The WebSocket configuration containing the configuration for this WebSocket instance.
856    ///
857    /// # Panics
858    ///
859    /// Panics if the context in the WebSocket configuration is not set (i.e., it's the default context).
860    /// Panics if the broadcast type in the WebSocket configuration is `BroadcastType::Unknown`.
861    pub async fn run<B>(&self, websocket_config: WebSocketConfig<'_, B>)
862    where
863        B: BroadcastTypeTrait,
864    {
865        let capacity: Capacity = websocket_config.get_capacity();
866        let broadcast_type: BroadcastType<B> = websocket_config.get_broadcast_type().clone();
867        let connected_hook: ServerHookHandler = websocket_config.get_connected_hook().clone();
868        let sended_hook: ServerHookHandler = websocket_config.get_sended_hook().clone();
869        let request_hook: ServerHookHandler = websocket_config.get_request_hook().clone();
870        let closed_hook: ServerHookHandler = websocket_config.get_closed_hook().clone();
871        let WebSocketConfig {
872            stream,
873            context: ctx,
874            ..
875        } = websocket_config;
876        let mut receiver: Receiver<Vec<u8>> = match &broadcast_type {
877            BroadcastType::PointToPoint(key1, key2) => self.point_to_point(key1, key2, capacity),
878            BroadcastType::PointToGroup(key) => self.point_to_group(key, capacity),
879            BroadcastType::Unknown => panic!("BroadcastType must be PointToPoint or PointToGroup"),
880        };
881        let key: String = BroadcastType::get_key(broadcast_type);
882        if connected_hook(stream, ctx).await.is_reject() {
883            return;
884        }
885        let mut is_reject: bool;
886        loop {
887            tokio::select! {
888                request_res = stream.try_get_websocket_request() => {
889                    if let Ok(body) = request_res {
890                        ctx.get_mut_request().set_body(body);
891                        is_reject = request_hook(stream, ctx).await.is_reject();
892                    } else {
893                        is_reject = true;
894                        closed_hook(stream, ctx).await;
895                    }
896                    let body: ResponseBody = ctx.get_response().get_body().clone();
897                    let is_err: bool = self.get_broadcast_map().try_send(&key, body).is_err();
898                    if is_err || sended_hook(stream, ctx).await.is_reject() || is_reject {
899                        break;
900                    }
901                },
902                msg_res = receiver.recv() => {
903                    if let Ok(msg) = &msg_res {
904                        if stream.try_send_list(&WebSocketFrame::create_frame_list(msg)).await.is_ok() {
905                            continue;
906                        } else {
907                            break;
908                        }
909                    }
910                    break;
911                }
912            }
913        }
914        stream.set_closed(true);
915    }
916}