1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
//! Module handling WebSocket connections and relaying messages between clients and rooms.
use std::fmt::Debug;
use std::sync::Arc;
use std::time::Duration;
use anyhow::Result;
use axum::extract::ws::{Message, WebSocket};
use futures::{FutureExt, SinkExt, StreamExt};
use redis::AsyncCommands;
use serde::{Deserialize, Serialize};
use sysinfo::System;
use thiserror::Error;
use tokio::sync::{watch, Mutex, RwLock};
use tokio::time::sleep;
use tokio_tungstenite::connect_async;
use tokio_tungstenite::tungstenite::protocol::Message as TungsteniteMessage;
use tracing::{debug, error, info, trace};
use yrs::sync::Awareness;
use yrs::Doc;
use crate::broadcast::{BroadcastManager, BroadcastManagerError};
use crate::ids::IdFactory;
use crate::RedisKeygenerator;
/// Struct for serializing node information.
#[derive(Serialize)]
struct NodeInfo {
address: String,
num_rooms: usize,
num_connections: usize,
cpu_usage: f32,
total_memory: u64,
used_memory: u64,
}
/// Struct for serializing and deserializing room information.
#[derive(Serialize, Deserialize)]
struct RoomInfo {
address: String,
node_id: String,
participants: Option<usize>,
}
/// Represents a relay node responsible for handling client connections and room management.
#[derive(derive_builder::Builder)]
#[builder(build_fn(name = "build_inner", validate = "Self::validate"))]
pub struct RelayNode {
/// The network address of the relay node.
pub address: String,
/// The unique identifier of the relay node.
pub id: String,
/// Redis client for interacting with the Redis server.
#[builder(setter(into))]
redis: redis::Client,
/// Manages broadcasting to clients across rooms.
#[builder(default = "Arc::new(BroadcastManager::new())")]
broadcast_manager: Arc<BroadcastManager>,
/// Factory for generating unique IDs.
#[builder(default = "Arc::new(IdFactory::new())")]
id_factory: Arc<IdFactory>,
}
impl Debug for RelayNode {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("RelayNode")
.field("address", &self.address)
.field("id", &self.id)
.finish()
}
}
impl RelayNode {
/// Creates a new `RelayNode` and starts the node info worker.
///
/// # Arguments
///
/// * `address` - The network address of the relay node.
/// * `redis` - A Redis client.
pub fn new(address: String, redis: redis::Client) -> Self {
let factory = IdFactory::new();
let relay_node = RelayNode {
address: address.clone(),
id: factory.gen_id(),
redis: redis.clone(),
broadcast_manager: Arc::new(BroadcastManager::new()),
id_factory: Arc::new(factory),
};
// Start the node info worker
relay_node.start_node_info_worker();
relay_node
}
pub fn builder() -> RelayNodeBuilder {
RelayNodeBuilder::default()
}
/// Starts a background task that periodically reports node info to Redis.
#[tracing::instrument]
fn start_node_info_worker(&self) {
let redis_clone = self.redis.clone();
let broadcast_manager_clone = self.broadcast_manager.clone();
let address_clone = self.address.clone();
let id = self.id.clone();
// Spawn the node info worker as a background task
tokio::spawn(async move {
// Initialize system info collector
let mut sys = System::new_all();
loop {
// Update system info
sys.refresh_all();
// Collect data
let num_rooms = broadcast_manager_clone.num_rooms();
let num_connections = broadcast_manager_clone.total_listeners();
// Collect host information
let cpu_usage = sys.global_cpu_usage();
let total_memory = sys.total_memory();
let used_memory = sys.used_memory();
// Create the info to report
let node_info = NodeInfo {
address: address_clone.clone(),
num_rooms,
num_connections,
cpu_usage,
total_memory,
used_memory,
};
// Report to Redis
let mut redis_conn = match redis_clone.get_multiplexed_async_connection().await {
Ok(conn) => conn,
Err(err) => {
error!("Failed to get Redis connection: {}", err);
sleep(Duration::from_secs(5)).await;
continue;
}
};
let node_info_key = RedisKeygenerator::node_key(&id);
// Serialize node_info to JSON
let node_info_json = match serde_json::to_string(&node_info) {
Ok(json) => json,
Err(err) => {
error!("Failed to serialize node info: {}", err);
sleep(Duration::from_secs(5)).await;
continue;
}
};
// Store in Redis with an expiry
if let Err(err) = redis_conn
.set_ex::<_, _, ()>(&node_info_key, node_info_json, 10)
.await
{
error!("Failed to set node info in Redis: {}", err);
}
// Sleep before next report
sleep(Duration::from_secs(5)).await;
}
});
}
/// Starts a background task that periodically reports room info to Redis.
///
/// # Arguments
///
/// * `room_name` - The name of the room.
/// * `shutdown_rx` - A receiver for shutdown signals.
#[tracing::instrument]
fn start_room_info_worker(&self, room_name: String, mut shutdown_rx: watch::Receiver<()>) {
let redis_clone = self.redis.clone();
let id = self.id.clone();
let address_clone = self.address.clone();
let broadcast_clone = self.broadcast_manager.clone();
let room_key = RedisKeygenerator::room_key(&room_name);
// Spawn the room info worker as a background task
tokio::spawn(async move {
loop {
// Check for shutdown signal
if shutdown_rx.has_changed().unwrap_or(false) {
info!("Shutting down room info worker for room '{}'", room_name);
break;
}
// Create the room info to report
let room_info = RoomInfo {
address: address_clone.clone(),
node_id: id.clone(),
participants: broadcast_clone.listeners(&room_name),
};
// Report to Redis
let mut redis_conn = match redis_clone.get_multiplexed_async_connection().await {
Ok(conn) => conn,
Err(err) => {
error!("Failed to get Redis connection: {}", err);
sleep(Duration::from_secs(5)).await;
continue;
}
};
// Serialize room_info to JSON
let room_info_json = match serde_json::to_string(&room_info) {
Ok(json) => json,
Err(err) => {
error!("Failed to serialize room info: {}", err);
sleep(Duration::from_secs(5)).await;
continue;
}
};
// Store in Redis with an expiry (TTL)
if let Err(err) = redis_conn
.set_ex::<_, _, ()>(&room_key, room_info_json, 10)
.await
{
error!("Failed to set room info in Redis: {}", err);
}
trace!("Room worker updated room info for '{}'", room_name);
// Sleep before next report
tokio::select! {
_ = sleep(Duration::from_secs(5)) => {},
_ = shutdown_rx.changed() => {
info!("Shutting down room info worker for room '{}'", room_name);
break;
}
}
}
});
}
/// Handles a WebSocket upgrade by routing the connection to the appropriate room or building a relay.
///
/// # Arguments
///
/// * `socket` - The WebSocket connection.
/// * `room_name` - The name of the room to connect to.
#[tracing::instrument(skip(socket))]
pub async fn handle_upgrade(&self, socket: WebSocket, room_name: String) {
let room_key = RedisKeygenerator::room_key(&room_name);
let mut redis_conn = match self.redis.get_multiplexed_async_connection().await {
Ok(conn) => conn,
Err(err) => {
error!("Failed to get Redis connection: {}", err);
return;
}
};
// Attempt to get the room info from Redis
let room_info_json: Option<String> = match redis_conn.get(&room_key).await {
Ok(info) => info,
Err(err) => {
error!("Failed to get room info from Redis: {}", err);
return;
}
};
if let Some(room_info_json) = room_info_json {
// Parse the room info
let room_info: RoomInfo = match serde_json::from_str(&room_info_json) {
Ok(info) => info,
Err(err) => {
error!("Failed to parse room info JSON: {}", err);
return;
}
};
if room_info.address == self.address {
// Room is hosted on this server
debug!("Handling connection to '{}' locally", room_name);
self.handle_socket(socket, room_name.clone()).await;
} else {
debug!(
"Expecting room '{}' to be hosted at {}, building relay",
room_name, room_info.address
);
// Room is on a different server; build a relay
self.build_relay(socket, room_info.address.clone(), room_name.clone())
.await;
}
} else {
// Room does not exist or TTL expired; attempt to create it atomically
// Create room locally so that once new connections get sent the room is ready
match self
.broadcast_manager
.create_room(
&room_name,
Arc::new(RwLock::new(Awareness::new(Doc::new()))),
)
.await
{
Ok(_) => {}
Err(e) => {
error!(
"Failed to create room '{}' in broadcast manager: {}",
room_name,
e.to_string()
);
return;
}
}
// Try to set the room info in Redis with NX and TTL
let room_info = RoomInfo {
address: self.address.clone(),
node_id: self.id.clone(),
participants: self.broadcast_manager.listeners(&room_name),
};
// Serialize room_info to JSON
let room_info_json = match serde_json::to_string(&room_info) {
Ok(json) => json,
Err(err) => {
error!("Failed to serialize room info: {}", err);
return;
}
};
// Use a Lua script to set the room info with NX and TTL atomically
let script = redis::Script::new(
r#"
local setnx = redis.call('SETNX', KEYS[1], ARGV[1])
if setnx == 1 then
redis.call('EXPIRE', KEYS[1], tonumber(ARGV[2]))
end
return setnx
"#,
);
let ttl_seconds = 10;
let set_result: i32 = match script
.key(&room_key)
.arg(&room_info_json)
.arg(ttl_seconds)
.invoke_async(&mut redis_conn)
.await
{
Ok(result) => result,
Err(err) => {
error!("Failed to set room info in Redis: {}", err);
return;
}
};
if set_result == 1 {
// Successfully created the room in Redis
debug!("Room '{}' didn't exist so I created it", room_name);
// Start the room info worker to refresh TTL
let (shutdown_tx, shutdown_rx) = watch::channel(());
self.broadcast_manager
.store_room_shutdown_signal(&room_name, shutdown_tx)
.await;
self.start_room_info_worker(room_name.clone(), shutdown_rx);
// Handle the socket
self.handle_socket(socket, room_name.clone()).await;
} else {
// Another server created the room; get the updated room info
// Tear down local room
if let Err(e) = self.broadcast_manager.drop_room(&room_name) {
error!(
"Failed to remove pre-created room from broadcast manager: {}",
e
);
}
// Get new room info
let room_info_json: String = match redis_conn.get(&room_key).await {
Ok(info) => info,
Err(err) => {
error!("Failed to get room info from Redis after set_nx: {}", err);
return;
}
};
let room_info: RoomInfo = match serde_json::from_str(&room_info_json) {
Ok(info) => info,
Err(err) => {
error!("Failed to parse room info JSON: {}", err);
return;
}
};
debug!(
"Tried to build room '{}' but {} got to it first, building relay",
room_name, room_info.address
);
self.build_relay(socket, room_info.address.clone(), room_name.clone())
.await;
}
}
}
/// Builds a relay between the client and the current room server.
///
/// # Arguments
///
/// * `socket` - The WebSocket connection from the client.
/// * `current_room_server` - The address of the current room server.
/// * `room_name` - The name of the room.
#[tracing::instrument(skip(socket))]
pub async fn build_relay(
&self,
socket: WebSocket,
current_room_server: String,
room_name: String,
) {
let mut current_server_address = current_room_server;
let mut socket = Some(socket);
loop {
let room_server_url = format!("ws://{}/{}", current_server_address, room_name);
match connect_async(&room_server_url).await {
Ok((room_ws_stream, _)) => {
// Split both sockets
let (mut client_ws_sink, mut client_ws_stream) = socket.take().unwrap().split();
let (mut room_ws_sink, mut room_ws_stream) = room_ws_stream.split();
// Relay messages from client to room server
let client_to_room = async {
while let Some(msg) = client_ws_stream.next().await {
if let Ok(msg) = msg {
let msg = match msg {
Message::Text(text) => TungsteniteMessage::Text(text),
Message::Binary(bin) => TungsteniteMessage::Binary(bin),
_ => continue,
};
trace!("Passing message {} to {}", msg, current_server_address);
if room_ws_sink.send(msg).await.is_err() {
break;
}
} else {
break;
}
}
};
// Relay messages from room server to client
let room_to_client = async {
while let Some(msg) = room_ws_stream.next().await {
if let Ok(msg) = msg {
let msg = match msg {
TungsteniteMessage::Text(text) => Message::Text(text),
TungsteniteMessage::Binary(bin) => Message::Binary(bin),
_ => continue,
};
if client_ws_sink.send(msg).await.is_err() {
break;
}
} else {
break;
}
}
};
// Run both futures concurrently
futures::future::select(client_to_room.boxed(), room_to_client.boxed()).await;
break; // Exit the loop after successful relay
}
Err(err) => {
error!(
"Failed to connect to room server at {}: {}",
room_server_url, err
);
// Assume the room server is dead; attempt to take over the room
if let Some(s) = socket.take() {
let result = self.attempt_room_takeover(s, room_name.clone()).await;
match result {
Ok(new_socket) => {
// If takeover was successful, handle the socket directly
self.handle_socket(new_socket, room_name.clone()).await;
break;
}
Err(Some((new_socket, new_server_address))) => {
// Another server took over; retry with the new server address
socket = Some(new_socket);
current_server_address = new_server_address;
continue; // Retry the loop with the new server address
}
Err(None) => {
// Failed to take over and no new server address; cannot proceed
break;
}
}
} else {
// Socket has already been taken; cannot proceed
break;
}
}
}
}
}
/// Attempts to take over a room if the current server is unresponsive.
///
/// # Arguments
///
/// * `socket` - The WebSocket connection from the client.
/// * `room_name` - The name of the room.
///
/// # Returns
///
/// A `Result` indicating success or an `Option` containing the socket and new server address.
#[tracing::instrument(skip(socket))]
async fn attempt_room_takeover(
&self,
socket: WebSocket,
room_name: String,
) -> Result<WebSocket, Option<(WebSocket, String)>> {
info!("Attempting takeover of room {}", room_name);
let room_key = RedisKeygenerator::room_key(&room_name);
let mut redis_conn = match self.redis.get_multiplexed_async_connection().await {
Ok(conn) => conn,
Err(err) => {
error!("Failed to get Redis connection: {}", err);
return Err(None);
}
};
// Get the current room info and TTL
let room_info_json: Option<String> = match redis_conn.get(&room_key).await {
Ok(info) => info,
Err(err) => {
error!("Failed to get room info from Redis: {}", err);
return Err(None);
}
};
let ttl: i32 = match redis_conn.ttl(&room_key).await {
Ok(ttl) => ttl,
Err(err) => {
error!("Failed to get TTL from Redis: {}", err);
return Err(None);
}
};
// If the room info does not exist or TTL is expired, attempt to take over
let can_takeover = room_info_json.is_none() || ttl <= 0;
if can_takeover {
// Try to set the room info in Redis with NX and TTL
let room_info = RoomInfo {
address: self.address.clone(),
node_id: self.id.clone(),
participants: self.broadcast_manager.listeners(&room_name),
};
// Serialize room_info to JSON
let room_info_json = match serde_json::to_string(&room_info) {
Ok(json) => json,
Err(err) => {
error!("Failed to serialize room info: {}", err);
return Err(None);
}
};
// Use a Lua script to set the room info with NX and TTL atomically
let script = redis::Script::new(
r#"
local setnx = redis.call('SETNX', KEYS[1], ARGV[1])
if setnx == 1 then
redis.call('EXPIRE', KEYS[1], tonumber(ARGV[2]))
end
return setnx
"#,
);
let ttl_seconds = 10;
let set_result: i32 = match script
.key(room_key.clone())
.arg(&room_info_json)
.arg(ttl_seconds)
.invoke_async(&mut redis_conn)
.await
{
Ok(result) => result,
Err(err) => {
error!("Failed to set room info in Redis: {}", err);
return Err(None);
}
};
if set_result == 1 {
// Successfully took over the room
// Start the room info worker to refresh TTL
let (shutdown_tx, shutdown_rx) = watch::channel(());
self.broadcast_manager
.store_room_shutdown_signal(&room_name, shutdown_tx)
.await;
self.start_room_info_worker(room_name.clone(), shutdown_rx);
Ok(socket)
} else {
// Another server took over the room; get the new server address
let room_info_json: String = match redis_conn.get(&room_key).await {
Ok(info) => info,
Err(err) => {
error!("Failed to get room info from Redis after set_nx: {}", err);
return Err(None);
}
};
let room_info: RoomInfo = match serde_json::from_str(&room_info_json) {
Ok(info) => info,
Err(err) => {
error!("Failed to parse room info JSON: {}", err);
return Err(None);
}
};
Err(Some((socket, room_info.address)))
}
} else {
// Another process updated the key; get the new server address
let room_info_json = room_info_json.unwrap();
let room_info: RoomInfo = match serde_json::from_str(&room_info_json) {
Ok(info) => info,
Err(err) => {
error!(
"Failed to parse room info JSON from Redis after DEL failed: {}",
err
);
return Err(None);
}
};
Err(Some((socket, room_info.address)))
}
}
/// Handles the WebSocket connection for a specific room.
///
/// # Arguments
///
/// * `socket` - The WebSocket connection.
/// * `room_name` - The name of the room.
#[tracing::instrument(skip(socket))]
pub async fn handle_socket(&self, socket: WebSocket, room_name: String) {
let client_id = Arc::new(self.id_factory.gen_id());
info!("Client {} connected to room '{}'", client_id, room_name);
let room_key = RedisKeygenerator::room_key(&room_name);
let mut redis_conn = match self.redis.get_multiplexed_async_connection().await {
Ok(conn) => conn,
Err(err) => {
error!("Failed to get Redis connection: {}", err);
return;
}
};
// Split the WebSocket into sender and receiver
let (ws_tx, ws_rx) = socket.split();
// Prepare the sink to send binary messages
let sink = ws_tx.with(|data: Vec<u8>| {
futures::future::ok::<Message, axum::Error>(Message::Binary(data))
});
let sink = Arc::new(Mutex::new(sink));
// Clone Arc to use inside the closure
let client_id_clone = Arc::clone(&client_id);
// Process incoming messages
let stream = ws_rx
.map(move |result| {
let client_id = Arc::clone(&client_id_clone);
result
.map_err(WebSocketError::ReceiveError)
.and_then(|msg| match msg {
Message::Binary(data) => {
trace!("Client {} sent binary data: {:?}", client_id, data);
Ok(data)
}
Message::Text(text) => {
trace!("Client {} sent text data: {}", client_id, text);
Ok(text.into_bytes())
}
_ => {
// Unsupported message type, return error
Err(WebSocketError::UnsupportedMessageType)
}
})
})
.take_while(|res| futures::future::ready(res.is_ok())); // Stop on first error
// Subscribe the client to the room
let subscription = match self
.broadcast_manager
.subscribe(&room_name, sink.clone(), stream)
{
Ok(sub) => sub,
Err(e) => {
error!("Client {} subscription error: {}", client_id, e);
// Close the WebSocket connection
if let Err(e) = sink.lock().await.close().await {
error!("Failed to close WebSocket: {}", e);
}
return;
}
};
// Wait for the subscription to complete or be cancelled
if let Err(e) = subscription.completed().await {
error!("Client {} subscription error: {}", client_id, e);
}
// Close the WebSocket connection
if let Err(e) = sink.lock().await.close().await {
error!("Failed to close WebSocket: {}", e);
}
// Unsubscribe the client when the connection is closed
if let Err(e) = self.broadcast_manager.unsubscribe(&room_name) {
error!("Client {} unsubscription error: {}", client_id, e);
}
// After unsubscribing, get the updated listener count
if let Some(0) = self.broadcast_manager.listeners(&room_name) {
// Attempt to drop the room
match self.broadcast_manager.drop_room(&room_name) {
Ok(_) => {
// Room was successfully dropped
debug!("Room '{}' dropped", room_name);
// Send shutdown signal to the room info worker
if let Some(shutdown_tx) = self
.broadcast_manager
.remove_room_shutdown_signal(&room_name)
.await
{
let _ = shutdown_tx.send(());
}
// Remove room info from Redis
match redis_conn.del::<&String, i32>(&room_key).await {
Ok(_) => {}
Err(e) => {
error!("Failed to drop room_key '{}' from Redis: {}", room_key, e);
}
}
}
Err(BroadcastManagerError::StillParticipants { .. }) => {
// Room still has participants; do not proceed
debug!("Room '{}' still has participants; not dropping", room_name);
}
Err(e) => {
error!("Failed to drop room '{}': {}", room_name, e);
}
}
} else {
// Room still has listeners or does not exist; do nothing
debug!(
"Room '{}' not dropped. Listeners: {:?}",
room_name,
self.broadcast_manager.listeners(&room_name)
);
}
info!(
"Client {} disconnected from room '{}'",
client_id, room_name
);
}
}
/// Errors that can occur when handling WebSocket messages.
#[derive(Debug, Error)]
enum WebSocketError {
/// An error occurred while receiving a WebSocket message.
#[error("WebSocket receive error: {0}")]
ReceiveError(#[from] axum::Error),
/// An unsupported message type was received.
#[error("Unsupported message type")]
UnsupportedMessageType,
}
unsafe impl Send for WebSocketError {}
unsafe impl Sync for WebSocketError {}
impl RelayNodeBuilder {
pub fn build(&self) -> Result<RelayNode, RelayNodeBuilderError> {
// Ensure that the required `redis` field is set
let redis = self
.redis
.clone()
.ok_or(RelayNodeBuilderError::UninitializedField("redis"))?;
// Ensure that `address` is set
let address = self
.address
.clone()
.ok_or(RelayNodeBuilderError::UninitializedField("address"))?;
// Handle `id`: use provided or generate using `id_factory`
let id = match &self.id {
Some(id_val) => id_val.clone(),
None => {
let id_factory = self
.id_factory
.clone()
.unwrap_or_else(|| Arc::new(IdFactory::new()));
id_factory.gen_id()
}
};
// Handle `broadcast_manager`
let broadcast_manager = self
.broadcast_manager
.clone()
.unwrap_or_else(|| Arc::new(BroadcastManager::new()));
// Handle `id_factory`
let id_factory = self
.id_factory
.clone()
.unwrap_or_else(|| Arc::new(IdFactory::new()));
let node = RelayNode {
address,
id,
redis,
broadcast_manager,
id_factory,
};
node.start_node_info_worker();
Ok(node)
}
pub fn validate(&self) -> Result<(), String> {
if let None = self.redis {
return Err("You must initialize a RelayNode with a redis client".to_string());
}
Ok(())
}
}