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, parse_config::FileMetadata, 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    /// Whether we have already requested and received headers (only need them once)
42    pub headers_received: bool,
43}
44
45impl WsConnection {
46    /// Create a new WsConnection
47    pub fn new(ws_sender: WsSender, ws_receiver: WsReceiver) -> Self {
48        log::info!("🔌 WsConnection::new() - creating connection (connecting=true, available=false)");
49        Self {
50            ws_sender,
51            ws_receiver,
52            msg_id: 1,
53            connecting: true,
54            asking_size_max: 1024,
55            available: false, // Start as unavailable until WsEvent::Opened is received
56            name: "".to_string(),
57            auto_loading: true, // Auto-loading enabled by default
58            waiting_for_response: false,
59            headers_received: false,
60        }
61    }
62
63    /// Connect to a WebSocket
64    /// # Errors
65    /// Return an error as String if the connection failed - see [`ewebsock::connect_with_wakeup`] for more details
66    pub fn connect(url: &str, wakeup: impl Fn() + Send + Sync + 'static) -> Result<(WsSender, WsReceiver), String> {
67        let options = ewebsock::Options {
68            #[cfg(not(target_arch = "wasm32"))]
69            additional_headers: vec![("Origin".to_string(), "tramex".to_string())],
70            ..Default::default()
71        };
72        ewebsock::connect_with_wakeup(url, options, wakeup)
73    }
74
75    /// Try to close the ws
76    /// # Errors
77    /// Return an error if its fail see [`ewebsock::WsSender::close`] for more details
78    pub fn close_impl(&mut self) -> Result<(), TramexError> {
79        self.ws_sender.close();
80        Ok(())
81    }
82}
83
84impl InterfaceTrait for WsConnection {
85    fn get_more_data(&mut self, layer_list: Layers, _data: &mut Data) -> Result<(), Vec<TramexError>> {
86        // Don't send request if already waiting for a response
87        if self.waiting_for_response {
88            return Ok(());
89        }
90
91        let msg = LogGet::new(self.msg_id, layer_list, self.asking_size_max);
92        // Request headers on the first request to get FileMetadata (technology, pci, etc.)
93        let msg = if !self.headers_received { msg.with_headers() } else { msg };
94        log::debug!("📤 Sending log_get request #{}", self.msg_id);
95        match serde_json::to_string(&msg) {
96            Ok(msg_stringed) => {
97                log::debug!("📤 JSON payload: {}", msg_stringed);
98                self.ws_sender.send(WsMessage::Text(msg_stringed));
99                self.msg_id += 1;
100                self.waiting_for_response = true;
101                log::debug!("⏳ Waiting for response...");
102            }
103            Err(err) => {
104                log::error!("Error encoding message: {err:?}");
105                return Err(vec![tramex_error!(
106                    err.to_string(),
107                    crate::errors::ErrorCode::WebSocketErrorEncodingMessage
108                )]);
109            }
110        }
111        Ok(())
112    }
113
114    fn close(&mut self) -> Result<(), TramexError> {
115        self.close_impl()
116    }
117
118    fn supports_preloading(&self) -> bool {
119        false // WebSocket cannot preload - server controls data
120    }
121
122    fn get_total_event_count(&self) -> Option<usize> {
123        None // Unknown for WebSocket
124    }
125
126    fn is_fully_read(&self) -> bool {
127        !self.available // If connection is closed, we're done
128    }
129}
130
131impl WsConnection {
132    /// Check if we should send a new request
133    /// Returns true if auto-loading is enabled and not currently waiting for a response
134    /// The server-side timeout and allow_empty=false handle the timing automatically
135    pub fn should_request_more(&self) -> bool {
136        self.auto_loading && !self.waiting_for_response
137    }
138
139    /// Try to receive data
140    /// # Errors
141    /// Return an error if the data is not received correctly
142    pub fn try_recv(&mut self, data: &mut Data) -> Result<(), Vec<TramexError>> {
143        while let Some(event) = self.ws_receiver.try_recv() {
144            log::debug!(
145                "🔵 WsConnection::try_recv() - event received, current state: connecting={}, available={}",
146                self.connecting,
147                self.available
148            );
149            match event {
150                WsEvent::Message(msg) => {
151                    // Only set available=true if we've already opened (connecting=false)
152                    // This handles the case where messages arrive before the Opened event
153                    if !self.connecting {
154                        self.available = true;
155                    }
156                    match msg {
157                        WsMessage::Text(event_text) => {
158                            // log::debug!("📨 Raw WebSocket text message: {}", event_text);
159                            let decoded: Result<WebSocketLog, serde_json::Error> = serde_json::from_str(&event_text);
160                            match decoded {
161                                Ok(decoded_data) => {
162                                    log::debug!("✅ Received log_get response with {} logs", decoded_data.logs.len());
163
164                                    // Mark that we received the response - ready for next request
165                                    self.waiting_for_response = false;
166
167                                    // Extract FileMetadata from top-level headers if present (first request only)
168                                    if !self.headers_received
169                                        && let Some(ref headers) = decoded_data.headers
170                                    {
171                                        data.metadata = FileMetadata::parse_from_lines(headers);
172                                        log::info!(
173                                            "📋 Parsed WebSocket metadata: technology={:?}, pci={:?}, mode={:?}",
174                                            data.metadata.technology,
175                                            data.metadata.pci,
176                                            data.metadata.mode
177                                        );
178                                        self.headers_received = true;
179                                    }
180
181                                    let mut errors = vec![];
182                                    for one_log in decoded_data.logs {
183                                        // Skip NR band combinations logs
184                                        if Layer::RRC == one_log.layer
185                                            && one_log.data.iter().any(|line| line.contains("NR band combinations"))
186                                        {
187                                            log::debug!("Skipping NR band combinations log");
188                                            continue;
189                                        }
190                                        match one_log.extract_data() {
191                                            Ok(trace) => {
192                                                data.events.push(trace);
193                                            }
194                                            Err(err) => {
195                                                log::error!("Error while extracting data: {err:?}");
196                                                errors.push(err);
197                                            }
198                                        }
199                                    }
200                                    if !errors.is_empty() {
201                                        return Err(errors);
202                                    }
203                                }
204                                Err(_err) => {
205                                    let decoded_base: Result<BaseMessage, serde_json::Error> =
206                                        serde_json::from_str(&event_text);
207                                    match decoded_base {
208                                        Ok(decoded_data) => {
209                                            if decoded_data.message == "ready" {
210                                                log::debug!(
211                                                    "✅ Received 'ready' message from server: {}",
212                                                    decoded_data.name
213                                                );
214                                                log::debug!(
215                                                    "💡 Server is ready. You need to click 'Load More' or enable auto-loading to request logs."
216                                                );
217                                            }
218                                            log::debug!("📥 Received BaseMessage: {decoded_data:?}");
219                                            self.name = decoded_data.name;
220                                        }
221                                        Err(err) => {
222                                            log::error!("Error decoding message: {err:?}");
223                                            log::error!("Message: {event_text:?}");
224                                            return Err(vec![tramex_error!(
225                                                err.to_string(),
226                                                crate::errors::ErrorCode::WebSocketErrorDecodingMessage
227                                            )]);
228                                        }
229                                    }
230                                }
231                            }
232                        }
233                        WsMessage::Unknown(str_error) => {
234                            log::error!("Unknown message: {str_error:?}");
235                            return Err(vec![tramex_error!(
236                                str_error,
237                                crate::errors::ErrorCode::WebSocketUnknownMessageReceived
238                            )]);
239                        }
240                        WsMessage::Binary(bin) => {
241                            log::error!("Unknown binary message: {bin:?}");
242                            return Err(vec![tramex_error!(
243                                format!("Unknown binary message: {bin:?}"),
244                                crate::errors::ErrorCode::WebSocketUnknownBinaryMessageReceived
245                            )]);
246                        }
247                        _ => {
248                            log::debug!("Received Ping-Pong")
249                        }
250                    }
251                }
252                WsEvent::Opened => {
253                    log::info!("✅ WsEvent::Opened - connection established! Setting connecting=false, available=true");
254                    self.connecting = false;
255                    self.available = true;
256                }
257                WsEvent::Closed => {
258                    log::warn!("⚠️ WsEvent::Closed - connection closed. Setting connecting=false, available=false");
259                    self.available = false;
260                    self.connecting = false;
261                    return Err(vec![tramex_error!(
262                        "WebSocket closed".to_string(),
263                        crate::errors::ErrorCode::WebSocketClosed
264                    )]);
265                }
266                WsEvent::Error(str_err) => {
267                    log::error!(
268                        "❌ WsEvent::Error - connection failed: {}. Setting connecting=false, available=false",
269                        str_err
270                    );
271                    self.available = false;
272                    self.connecting = false;
273                    return Err(vec![tramex_error!(str_err, crate::errors::ErrorCode::WebSocketError)]);
274                }
275            }
276        }
277        Ok(())
278    }
279}
280
281impl Debug for WsConnection {
282    fn fmt(&self, formatter: &mut Formatter<'_>) -> core::fmt::Result {
283        formatter
284            .debug_struct("Interface")
285            .field("ws_sender", &"Box<WsSender>")
286            .field("ws_receiver", &"Box<WsReceiver>")
287            .field("connecting", &self.connecting)
288            .finish()
289    }
290}
291
292impl Drop for WsConnection {
293    fn drop(&mut self) {
294        log::debug!("Cleaning WsConnection");
295        self.ws_sender.close()
296    }
297}