Skip to main content

injoint_client/joint_client/
mod.rs

1use crate::event_listener::{EventListener, Handler};
2use crate::message::MessageRequest;
3use crate::utils::{WSSink, WSStream};
4use anyhow::{anyhow, Result};
5use futures_util::{SinkExt, StreamExt};
6use serde::de::DeserializeOwned;
7use serde::Serialize;
8use std::sync::Arc;
9use tokio::sync::Mutex;
10use tokio_tungstenite::{connect_async, tungstenite::protocol::Message};
11
12pub struct JointClient<A, T>
13where
14    A: DeserializeOwned + Serialize + Send + Sync + 'static,
15    T: DeserializeOwned + Send + Sync + 'static,
16{
17    listener: Arc<Mutex<EventListener<A, T>>>,
18    pub client_id: Option<Arc<Mutex<u64>>>,
19    pub room_id: Option<Arc<Mutex<u64>>>,
20    sink: Option<Arc<Mutex<WSSink>>>,
21    stream: Option<Arc<Mutex<WSStream>>>,
22}
23
24impl<A, T> JointClient<A, T>
25where
26    A: DeserializeOwned + Serialize + Send + Sync + 'static,
27    T: DeserializeOwned + Send + Sync + 'static,
28{
29    pub fn new() -> Self {
30        JointClient {
31            listener: Arc::new(Mutex::new(EventListener::new())),
32            client_id: None,
33            room_id: None,
34            sink: None,
35            stream: None,
36        }
37    }
38
39    pub async fn register_handler(&mut self, handler: Handler<A, T>) {
40        let listener = self.listener.lock().await;
41        listener.register_handler(handler).await
42    }
43
44    pub async fn connect(&mut self, addr: &str) -> Result<()> {
45        let (ws_stream, _) = connect_async(addr).await?;
46        let (sink, stream) = ws_stream.split();
47        self.sink = Some(Arc::new(Mutex::new(sink)));
48        self.stream = Some(Arc::new(Mutex::new(stream)));
49        Ok(())
50    }
51
52    pub async fn listen(&self) -> Result<()> {
53        if let Some(stream) = self.stream.clone() {
54            let listener = self.listener.clone();
55
56            tokio::spawn(async move {
57                let mut stream = stream.lock().await;
58                while let Some(msg) = stream.next().await {
59                    match msg {
60                        Ok(Message::Text(text)) => {
61                            let listener = listener.lock().await; // Ensure listener is locked within the async block
62                            if let Err(e) = listener.handle_event(&text).await {
63                                eprintln!("Error handling message: {}", e);
64                            }
65                        }
66                        Ok(_) => {}
67                        Err(e) => {
68                            eprintln!("WebSocket error: {}", e);
69                            break;
70                        }
71                    }
72                }
73            });
74
75            Ok(())
76        } else {
77            Err(anyhow!("WebSocket stream is not initialized"))
78        }
79    }
80
81    pub async fn create_room(&self) -> Result<()> {
82        self.send_event(MessageRequest::create_room()).await
83    }
84
85    pub async fn join_room(&self, room_id: u64) -> Result<()> {
86        self.send_event(MessageRequest::join_room(room_id)).await
87    }
88
89    pub async fn leave_room(&self) -> Result<()> {
90        self.send_event(MessageRequest::leave_room()).await
91    }
92
93    pub async fn dispatch_action(&self, action: A) -> Result<()> {
94        self.send_event(MessageRequest::action(action)).await
95    }
96
97    async fn send_event(&self, request: MessageRequest<A>) -> Result<()> {
98        let create_json = serde_json::to_string(&request)?;
99        if let Some(sink) = &self.sink {
100            let mut sink = sink.lock().await;
101            sink.send(Message::Text(create_json.into())).await?;
102            Ok(())
103        } else {
104            Err(anyhow!("WebSocket stream is not initialized"))
105        }
106    }
107}