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}