tramex_tools/interface/websocket/
ws_connection.rs

1//! WsConnection struct
2use core::fmt::{Debug, Formatter};
3use ewebsock::{WsEvent, WsMessage, WsReceiver, WsSender};
4use std::vec;
5
6use crate::interface::interface_types::InterfaceTrait;
7use crate::interface::types::BaseMessage;
8use crate::tramex_error;
9use crate::{data::Data, errors::TramexError};
10
11use crate::interface::{layer::Layer, layer::Layers, log_get::LogGet, types::WebSocketLog};
12/// WsConnection struct
13pub struct WsConnection {
14    /// WebSocket sender
15    pub ws_sender: WsSender,
16
17    /// WebSocket receiver
18    pub ws_receiver: WsReceiver,
19
20    /// Message ID
21    pub msg_id: u64,
22
23    /// Connecting flag
24    pub connecting: bool,
25
26    /// Asking size max
27    pub asking_size_max: u64,
28
29    /// Available flag
30    pub available: bool,
31
32    /// Name of the receiver
33    pub name: String,
34
35    /// Auto-loading flag - when true, automatically requests more data
36    pub auto_loading: bool,
37
38    /// Waiting for response flag - true when a request has been sent and we're waiting for the response
39    pub waiting_for_response: bool,
40}
41
42impl WsConnection {
43    /// Create a new WsConnection
44    pub fn new(ws_sender: WsSender, ws_receiver: WsReceiver) -> Self {
45        Self {
46            ws_sender,
47            ws_receiver,
48            msg_id: 1,
49            connecting: true,
50            asking_size_max: 1024,
51            available: true,
52            name: "".to_string(),
53            auto_loading: true, // Auto-loading enabled by default
54            waiting_for_response: false,
55        }
56    }
57
58    /// Connect to a WebSocket
59    /// # Errors
60    /// Return an error as String if the connection failed - see [`ewebsock::connect_with_wakeup`] for more details
61    pub fn connect(url: &str, wakeup: impl Fn() + Send + Sync + 'static) -> Result<(WsSender, WsReceiver), String> {
62        let options = ewebsock::Options {
63            #[cfg(not(target_arch = "wasm32"))]
64            additional_headers: vec![("Origin".to_string(), "tramex".to_string())],
65            ..Default::default()
66        };
67        ewebsock::connect_with_wakeup(url, options, wakeup)
68    }
69
70    /// Try to close the ws
71    /// # Errors
72    /// Return an error if its fail see [`ewebsock::WsSender::close`] for more details
73    pub fn close_impl(&mut self) -> Result<(), TramexError> {
74        self.ws_sender.close();
75        Ok(())
76    }
77}
78
79impl InterfaceTrait for WsConnection {
80    fn get_more_data(&mut self, layer_list: Layers, _data: &mut Data) -> Result<(), Vec<TramexError>> {
81        // Don't send request if already waiting for a response
82        if self.waiting_for_response {
83            return Ok(());
84        }
85
86        let msg = LogGet::new(self.msg_id, layer_list, self.asking_size_max);
87        log::debug!("📤 Sending log_get request #{}", self.msg_id);
88        match serde_json::to_string(&msg) {
89            Ok(msg_stringed) => {
90                log::debug!("📤 JSON payload: {}", msg_stringed);
91                self.ws_sender.send(WsMessage::Text(msg_stringed));
92                self.msg_id += 1;
93                self.waiting_for_response = true;
94                log::debug!("⏳ Waiting for response...");
95            }
96            Err(err) => {
97                log::error!("Error encoding message: {err:?}");
98                return Err(vec![tramex_error!(
99                    err.to_string(),
100                    crate::errors::ErrorCode::WebSocketErrorEncodingMessage
101                )]);
102            }
103        }
104        Ok(())
105    }
106
107    fn close(&mut self) -> Result<(), TramexError> {
108        self.close_impl()
109    }
110
111    fn supports_preloading(&self) -> bool {
112        false // WebSocket cannot preload - server controls data
113    }
114
115    fn get_total_event_count(&self) -> Option<usize> {
116        None // Unknown for WebSocket
117    }
118
119    fn is_fully_read(&self) -> bool {
120        !self.available // If connection is closed, we're done
121    }
122}
123
124impl WsConnection {
125    /// Check if we should send a new request
126    /// Returns true if auto-loading is enabled and not currently waiting for a response
127    /// The server-side timeout and allow_empty=false handle the timing automatically
128    pub fn should_request_more(&self) -> bool {
129        self.auto_loading && !self.waiting_for_response
130    }
131
132    /// Try to receive data
133    /// # Errors
134    /// Return an error if the data is not received correctly
135    pub fn try_recv(&mut self, data: &mut Data) -> Result<(), Vec<TramexError>> {
136        while let Some(event) = self.ws_receiver.try_recv() {
137            // log::debug!("🔵 WebSocket event received: {:?}", event);
138            log::debug!("🔵 WebSocket event received");
139            self.connecting = false;
140            match event {
141                WsEvent::Message(msg) => {
142                    self.available = true;
143                    match msg {
144                        WsMessage::Text(event_text) => {
145                            // log::debug!("📨 Raw WebSocket text message: {}", event_text);
146                            let decoded: Result<WebSocketLog, serde_json::Error> = serde_json::from_str(&event_text);
147                            match decoded {
148                                Ok(decoded_data) => {
149                                    log::debug!("✅ Received log_get response with {} logs", decoded_data.logs.len());
150
151                                    // Mark that we received the response - ready for next request
152                                    self.waiting_for_response = false;
153
154                                    let mut errors = vec![];
155                                    for one_log in decoded_data.logs {
156                                        // Skip NR band combinations logs
157                                        if Layer::RRC == one_log.layer
158                                            && one_log.data.iter().any(|line| line.contains("NR band combinations"))
159                                        {
160                                            log::debug!("Skipping NR band combinations log");
161                                            continue;
162                                        }
163                                        match one_log.extract_data() {
164                                            Ok(trace) => {
165                                                data.events.push(trace);
166                                            }
167                                            Err(err) => {
168                                                log::error!("Error while extracting data: {err:?}");
169                                                errors.push(err);
170                                            }
171                                        }
172                                    }
173                                    if !errors.is_empty() {
174                                        return Err(errors);
175                                    }
176                                }
177                                Err(_err) => {
178                                    let decoded_base: Result<BaseMessage, serde_json::Error> =
179                                        serde_json::from_str(&event_text);
180                                    match decoded_base {
181                                        Ok(decoded_data) => {
182                                            if decoded_data.message == "ready" {
183                                                log::debug!(
184                                                    "✅ Received 'ready' message from server: {}",
185                                                    decoded_data.name
186                                                );
187                                                log::debug!(
188                                                    "💡 Server is ready. You need to click 'Load More' or enable auto-loading to request logs."
189                                                );
190                                            }
191                                            log::debug!("📥 Received BaseMessage: {decoded_data:?}");
192                                            self.name = decoded_data.name;
193                                        }
194                                        Err(err) => {
195                                            log::error!("Error decoding message: {err:?}");
196                                            log::error!("Message: {event_text:?}");
197                                            return Err(vec![tramex_error!(
198                                                err.to_string(),
199                                                crate::errors::ErrorCode::WebSocketErrorDecodingMessage
200                                            )]);
201                                        }
202                                    }
203                                }
204                            }
205                        }
206                        WsMessage::Unknown(str_error) => {
207                            log::error!("Unknown message: {str_error:?}");
208                            return Err(vec![tramex_error!(
209                                str_error,
210                                crate::errors::ErrorCode::WebSocketUnknownMessageReceived
211                            )]);
212                        }
213                        WsMessage::Binary(bin) => {
214                            log::error!("Unknown binary message: {bin:?}");
215                            return Err(vec![tramex_error!(
216                                format!("Unknown binary message: {bin:?}"),
217                                crate::errors::ErrorCode::WebSocketUnknownBinaryMessageReceived
218                            )]);
219                        }
220                        _ => {
221                            log::debug!("Received Ping-Pong")
222                        }
223                    }
224                }
225                WsEvent::Opened => {
226                    self.available = true;
227                    log::debug!("✅ WebSocket connection opened successfully");
228                }
229                WsEvent::Closed => {
230                    self.available = false;
231                    log::debug!("WebSocket closed");
232                    return Err(vec![tramex_error!(
233                        "WebSocket closed".to_string(),
234                        crate::errors::ErrorCode::WebSocketClosed
235                    )]);
236                }
237                WsEvent::Error(str_err) => {
238                    self.available = false;
239                    log::error!("WebSocket error: {str_err:?}");
240                    return Err(vec![tramex_error!(str_err, crate::errors::ErrorCode::WebSocketError)]);
241                }
242            }
243        }
244        Ok(())
245    }
246}
247
248impl Debug for WsConnection {
249    fn fmt(&self, formatter: &mut Formatter<'_>) -> core::fmt::Result {
250        formatter
251            .debug_struct("Interface")
252            .field("ws_sender", &"Box<WsSender>")
253            .field("ws_receiver", &"Box<WsReceiver>")
254            .field("connecting", &self.connecting)
255            .finish()
256    }
257}
258
259impl Drop for WsConnection {
260    fn drop(&mut self) {
261        log::debug!("Cleaning WsConnection");
262        self.ws_sender.close()
263    }
264}