injoint_client/joint_client/
mod.rs1use 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; 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}